Powered by AppSignal & Oban Pro

Broadway Zarr Pipeline

docs/livebooks/broadway_pipeline.livemd

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

  1. Broadway provides supervision and retry semantics for chunk pipelines
  2. ExZarr.Broadway.ChunkProducer emits %Broadway.Message{data: {index, binary}}
  3. The producer is finite — it stops when chunks are exhausted (trap exits in Livebook)
  4. Use stream_chunks/2 directly for simpler pipelines without Broadway