Powered by AppSignal & Oban Pro

Processing 1TB Zarr Arrays

06_02_1tb_arrays.livemd

Processing 1TB Zarr Arrays

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

Pattern

For terabyte-scale arrays, combine streaming with Flow partitioning. Below is a small demo of the same shape as a multi-stage reduce.

alias ExZarr.Array
alias ExZarr.Gallery.Pack

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

vals = for i <- 0..(100 * 100 - 1), do: i * 1.0
:ok = Array.set_slice(array, Pack.pack(vals, :float64), start: {0, 0}, stop: {100, 100})

partials =
  array
  |> ExZarr.Flow.chunk_flow(stages: 2)
  |> Flow.map(fn {_idx, data} ->
    for <<val::float-little-64 <- data>>, reduce: 0.0 do
      acc -> acc + val
    end
  end)
  |> Enum.to_list()

%{partial_sums: partials, total: Enum.sum(partials)}

Use write_stream/3 with checkpoints for ingestion pipelines that may restart. See Cloud storage patterns for retry guidance.