Powered by AppSignal & Oban Pro

Processing 100GB Zarr Arrays

06_01_100gb_arrays.livemd

Processing 100GB Zarr Arrays

Mix.install([
  {:ex_zarr, "~> 1.2"}
  # {:ex_zarr, path: Path.join(__DIR__, "../../..")},
])

Problem

A large array exceeds available RAM. You need statistics or transforms without loading the full dataset.

Demo setup

We use a small in-memory array to show the streaming pattern. The same API scales to 100GB+ on disk or object storage.

alias ExZarr.Array
alias ExZarr.Gallery.Pack

{:ok, array} =
  Array.create(
    shape: {200, 200},
    chunks: {50, 50},
    dtype: :float64,
    compressor: :zlib,
    storage: :memory
  )

# Fill with a simple ramp
vals = for i <- 0..(200 * 200 - 1), do: i * 1.0
:ok = Array.set_slice(array, Pack.pack(vals, :float64), start: {0, 0}, stop: {200, 200})
:ok

Stream chunks with bounded concurrency

{sum, count} =
  array
  |> Array.stream_chunks(concurrency: 4, ordered: false)
  |> Enum.reduce({0.0, 0}, fn {_index, data}, {acc_sum, acc_count} ->
    chunk_sum =
      for <<val::float-little-64 <- data>>, reduce: 0.0 do
        acc -> acc + val
      end

    {acc_sum + chunk_sum, acc_count + div(byte_size(data), 8)}
  end)

mean = sum / count
%{sum: sum, count: count, mean: mean}

Memory budget

With 1MB chunks and concurrency 8, peak memory is roughly 8MB of chunk data plus decompression buffers. Keep concurrency below ~10% of available RAM divided by average chunk size.

When to use slices instead

For row-wise access (e.g. time series along dimension 0):

array
|> Array.stream_slices(0, concurrency: 2)
|> Enum.take(3)
|> Enum.map(fn {index, data} -> {index, byte_size(data)} end)