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.