Powered by AppSignal & Oban Pro

Distributed Processing (Experimental)

06_07_distributed.livemd

Distributed Processing (Experimental)

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

Pattern

Multi-node chunk processing is still an experimental stretch goal. Today, partition chunk indices across nodes manually.

alias ExZarr.Array
alias ExZarr.Gallery.Pack

{:ok, array} =
  Array.create(
    shape: {40, 40},
    chunks: {10, 10},
    dtype: :int32,
    compressor: :zlib,
    storage: :memory
  )

vals = Enum.to_list(1..(40 * 40))
:ok = Array.set_slice(array, Pack.pack(vals, :int32), start: {0, 0}, stop: {40, 40})

node_count = 4
node_index = 0

node_indices =
  array
  |> Array.stream_chunks(include_missing: true, metadata: true)
  |> Stream.map(& &1.index)
  |> Stream.filter(fn idx -> rem(elem(idx, 0), node_count) == node_index end)
  |> Enum.to_list()

%{node_index: node_index, assigned_chunks: length(node_indices), sample: Enum.take(node_indices, 5)}

Investigate Horde and PartitionSupervisor for future multi-node orchestration.