Broadway Zarr Pipeline
Introduction
This livebook demonstrates fault-tolerant image tile processing using ExZarr and Broadway.
Mix.install([
{:ex_zarr, "~> 1.2"},
# {:ex_zarr, path: Path.join(__DIR__, "../../..")}
{:broadway, "~> 1.0"},
{:kino, "~> 0.13"}
])
Create Sample Array
alias ExZarr.Array
{:ok, array} =
Array.create(
shape: {100, 100},
chunks: {10, 10},
dtype: :uint8,
storage: :memory
)
# Write gradient pattern
for x <- 0..9, y <- 0..9 do
data = for(_ <- 1..100, into: <<>>, do: <<x * 10 + y>>)
Array.set_slice(array, data, start: {x * 10, y * 10}, stop: {(x + 1) * 10, (y + 1) * 10})
end
Process Chunks with stream_chunks
results =
array
|> Array.stream_chunks(concurrency: 4, metadata: true)
|> Enum.map(fn %{index: {row, col}, data: tile} ->
%{
chunk_row: row,
chunk_col: col,
bytes: byte_size(tile),
sum: Enum.sum(:binary.bin_to_list(tile))
}
end)
|> Enum.sort_by(&{&1.chunk_row, &1.chunk_col})
Kino.DataTable.new(results)
Broadway Pipeline Pattern
ExZarr.Broadway.ChunkProducer is a finite producer: when every chunk has been
emitted it stops with :normal, and Broadway shuts the topology down. In Livebook,
Broadway.start_link/2 links to the evaluation process, so that shutdown would look
like “Evaluation process terminated - shutdown” unless we trap exits.
Message data is {chunk_index, binary} (see ExZarr.Broadway.ChunkProducer).
defmodule TilePipeline do
use Broadway
def start_link(array, collector) do
Broadway.start_link(__MODULE__,
name: __MODULE__,
producer: [
module:
{ExZarr.Broadway.ChunkProducer,
array: array, stream_opts: [concurrency: 1, include_missing: true]},
concurrency: 1
],
processors: [
default: [concurrency: 4, max_demand: 5, min_demand: 2]
],
context: %{collector: collector}
)
end
@impl Broadway
def handle_message(_processor, %Broadway.Message{data: {index, data}} = message, %{
collector: collector
}) do
sum = data |> :binary.bin_to_list() |> Enum.sum()
Agent.update(collector, fn acc -> [%{index: inspect(index), sum: sum} | acc] end)
Broadway.Message.put_data(message, {index, sum})
end
end
# Avoid linking Livebook's eval process to Broadway's shutdown
Process.flag(:trap_exit, true)
{:ok, collector} = Agent.start_link(fn -> [] end)
# Stop a leftover pipeline from a previous cell run
if pid = Process.whereis(TilePipeline), do: Broadway.stop(pid, :normal)
{:ok, broadway} = TilePipeline.start_link(array, collector)
# Wait for the finite producer to finish (topology exits :shutdown / :normal)
receive do
{:EXIT, ^broadway, reason} ->
IO.puts("Broadway finished: #{inspect(reason)}")
after
30_000 ->
Broadway.stop(broadway, :normal)
IO.puts("Broadway timed out; stopped manually")
end
results =
collector
|> Agent.get(& &1)
|> Enum.sort_by(& &1.index)
Agent.stop(collector)
Kino.DataTable.new(results)
Key Takeaways
- Broadway provides supervision and retry semantics for chunk pipelines
ExZarr.Broadway.ChunkProduceremits%Broadway.Message{data: {index, binary}}- The producer is finite — it stops when chunks are exhausted (trap exits in Livebook)
- Use
stream_chunks/2directly for simpler pipelines without Broadway