Powered by AppSignal & Oban Pro

Rheo GenStage, Flow & Broadway (v0.9.0)

notebooks/pipelines.livemd

Rheo GenStage, Flow & Broadway (v0.9.0)

Mix.install(
  [
    # From this repo: prefer the path dep. From HexDocs / standalone Livebook:
    {:rheo, "~> 0.9.0"},
    # {:rheo, path: Path.join(__DIR__, "..")},
    {:kino, "~> 0.14"},
    {:broadway, "~> 1.3"},
    {:flow, "~> 1.2"}
  ],
  config: [
    rheo: [
      start_on_application: false,
      clock: Rheo.Clock.System,
      default_lease_ms: 30_000,
      default_max_attempts: 5
    ]
  ]
)

Intro

Rheo answers: “Was this event delivered, leased, and settled for this group?”

GenStage, Flow, and Broadway answer: “How do we shape demand, parallelism, and batching in the BEAM?”

This notebook keeps that boundary sharp. You will feed the same %Rheo.Lease{} through:

  1. A plain GenStage consumer
  2. Flow map / reduce / window (with settlement rules)
  3. A Broadway pipeline

Rheo is not replaced by any of them — and none of them should invent a second delivery identity.

Sibling notebooks: Quickstart · Concepts · Backends · Index


Setup

We use ETS again, with two partitions and a handful of ticks. Each pipeline demo uses its own consumer group so they do not steal work from each other.

alias Rheo.{Lease, Producer}

case Process.whereis(Rheo) do
  nil -> :ok
  pid -> GenServer.stop(pid)
end

{:ok, _} = Rheo.start_link(backend: Rheo.Backend.ETS)
:ok = Rheo.ensure_indexes()

stream = "pipeline-events"
:ok = Rheo.create_stream(stream, partition_count: 2)
:ok = Rheo.create_group(stream, "genstage")
:ok = Rheo.create_group(stream, "broadway")

{:ok, _} =
  Rheo.append_batch(
    stream,
    for(i <- 1..12, do: %{type: "tick", n: i, key: "k-#{rem(i, 3)}"})
  )

:ok

Rule of thumb: pick one consumption surface per {rheo, stream, group} — either Rheo.Consumer or Rheo.Producer. Mixing both on the same group means competing fetchers.


1. GenStage: demand becomes leases

Rheo.Producer is a GenStage producer. Downstream demand turns into Rheo.fetch/3. Emitted events are leases, not bare payloads — because settlement still belongs to Rheo.

After you finish work:

  1. Settle durably (Rheo.ack/1 or Producer.ack/3)
  2. Free the producer's inflight slot (confirm/2 is included in Producer.ack/3)

Otherwise the producer keeps renewing the lease and will not fetch more once :max_demand is full.

defmodule Pipelines.Sink do
  @moduledoc "Tiny GenStage consumer that forwards leases to the Livebook process."
  use GenStage

  def start_link(opts), do: GenStage.start_link(__MODULE__, opts)

  @impl true
  def init(opts) do
    producer = Keyword.fetch!(opts, :producer)
    owner = Keyword.fetch!(opts, :owner)
    {:consumer, owner, subscribe_to: [{producer, max_demand: 4}]}
  end

  @impl true
  def handle_events(leases, _from, owner) do
    Enum.each(leases, &send(owner, {:lease, &1}))
    {:noreply, [], owner}
  end
end

{:ok, producer} =
  Producer.start_link(
    stream: stream,
    group: "genstage",
    max_demand: 4,
    poll_ms: 50
  )

{:ok, _sink} = Pipelines.Sink.start_link(producer: producer, owner: self())

genstage_leases =
  for _ <- 1..4 do
    receive do
      {:lease, %Lease{} = lease} -> lease
    after
      3_000 -> raise "timeout waiting for GenStage lease"
    end
  end

Enum.each(genstage_leases, fn lease ->
  # Prefer Producer.ack/3 in real code (settle + confirm together).
  :ok = Rheo.ack(lease)
  :ok = Producer.confirm(producer, lease.lease_id)
end)

%{
  received: length(genstage_leases),
  tick_numbers: Enum.map(genstage_leases, & &1.event.payload["n"]),
  inflight_after_confirm: Producer.inflight_count(producer)
}

inflight_after_confirm should be 0. Stop this producer before the Flow section so we do not leave a polling process around.

:ok = GenStage.stop(producer)

2. Flow: when is it safe to settle?

Flow is not a Rheo Hex package. You compose it yourself with Flow.from_stages/1. The lease type and settlement API stay the same as Broadway (and work against Redis the same way — see Backends).

Stage When to ACK
map (1:1) After the map side-effect succeeds
partition Still holding the lease; settle downstream
reduce After the aggregate is durable (on_trigger)
window After the window closes

Crashing before settle ⇒ lease expiry ⇒ at-least-once redelivery. Never ACK inside reduce “because Flow accepted the event”.

Livebook pitfall (read this)

Enum.take(n) only keeps n results. Upstream stages may still demand and run side effects for more events. If you ack inside Flow.map on a shared group, you can drain the whole group before the next cell runs — and the next cell appears to “hang forever”.

Each demo below therefore uses:

  • a fresh group
  • a finite number of appended events
  • a helper that times out and stops the producer
run_flow! = fn producer, flow, n, timeout ->
  task = Task.async(fn -> flow |> Enum.take(n) end)

  case Task.yield(task, timeout) || Task.shutdown(task, :brutal_kill) do
    {:ok, result} ->
      _ = GenStage.stop(producer)
      result

    nil ->
      _ = GenStage.stop(producer)
      raise "Flow did not finish in #{timeout}ms"
  end
end

:ok

Map 1:1 — settle per lease

:ok = Rheo.create_group(stream, "flow-map")

{:ok, _} =
  Rheo.append_batch(
    stream,
    for(i <- 1..4, do: %{type: "flow-map", n: i, key: "m-#{i}"})
  )

{:ok, flow_map_producer} =
  Producer.start_link(
    stream: stream,
    group: "flow-map",
    max_demand: 4,
    poll_ms: 50
  )

mapped =
  run_flow!.(
    flow_map_producer,
    Flow.from_stages([flow_map_producer])
    |> Flow.map(fn %Lease{} = lease ->
      :ok = Producer.ack(flow_map_producer, lease)
      {lease.event.payload["n"], lease.receipt}
    end),
    4,
    5_000
  )

%{mapped: mapped}

Reduce — settle only on trigger

We keep leases in the reducer accumulator and ACK them in on_trigger after a batch of three. That mimics “materialize the aggregate, then settle inputs”.

:ok = Rheo.create_group(stream, "flow-reduce")

{:ok, _} =
  Rheo.append_batch(
    stream,
    for(i <- 1..6, do: %{type: "flow-reduce", n: i, key: "r-#{rem(i, 2)}"})
  )

{:ok, flow_reduce_producer} =
  Producer.start_link(
    stream: stream,
    group: "flow-reduce",
    max_demand: 6,
    poll_ms: 50
  )

window = Flow.Window.global() |> Flow.Window.trigger_every(3)

[batch_size] =
  run_flow!.(
    flow_reduce_producer,
    Flow.from_stages([flow_reduce_producer])
    |> Flow.partition(key: fn _ -> :one end, window: window, stages: 1)
    |> Flow.reduce(fn -> [] end, fn %Lease{} = lease, acc -> [lease | acc] end)
    |> Flow.on_trigger(fn leases ->
      Enum.each(leases, fn lease ->
        :ok = Producer.ack(flow_reduce_producer, lease)
      end)

      {[length(leases)], []}
    end),
    1,
    5_000
  )

%{reduce_batch_size: batch_size}

reduce_batch_size should be 3 — the trigger size, not the full group.

Window — two flushes of two

Same idea with trigger_every(2) and Enum.take(2) so we observe both window emissions.

:ok = Rheo.create_group(stream, "flow-window")

{:ok, _} =
  Rheo.append_batch(
    stream,
    for(i <- 1..4, do: %{type: "flow-window", n: i, key: "same"})
  )

{:ok, flow_window_producer} =
  Producer.start_link(
    stream: stream,
    group: "flow-window",
    max_demand: 4,
    poll_ms: 50
  )

window = Flow.Window.global() |> Flow.Window.trigger_every(2)

batches =
  run_flow!.(
    flow_window_producer,
    Flow.from_stages([flow_window_producer])
    |> Flow.partition(key: fn _ -> :one end, window: window, stages: 1)
    |> Flow.reduce(fn -> [] end, fn lease, acc -> [lease | acc] end)
    |> Flow.on_trigger(fn leases ->
      Enum.each(leases, fn lease ->
        :ok = Producer.ack(flow_window_producer, lease)
      end)

      {[length(leases)], []}
    end),
    2,
    5_000
  )

%{window_batches: batches}

Expect window_batches == [2, 2].


3. Broadway: processors without owning delivery

Broadway is excellent at concurrency, batchers, and rate limits. Rheo remains the source of truth for leases.

Rheo.Broadway.transform/2 wraps each lease into a %Broadway.Message{}. Successful messages ACK; failed messages NACK (or reject if configured).

This uses the "broadway" group created in setup — it does not compete with GenStage/Flow groups above.

{:ok, collector} = Agent.start_link(fn -> [] end)

defmodule Pipelines.RiskBroadway do
  use Broadway

  def start_link(opts) do
    Broadway.start_link(__MODULE__,
      name: __MODULE__,
      context: %{collector: Keyword.fetch!(opts, :collector)},
      producer: [
        module:
          {Rheo.Producer,
           stream: Keyword.fetch!(opts, :stream),
           group: "broadway",
           max_demand: 20,
           poll_ms: 50},
        transformer: {Rheo.Broadway, :transform, []},
        concurrency: 1
      ],
      processors: [default: [concurrency: 4]]
    )
  end

  @impl true
  def handle_message(_processor, message, %{collector: collector}) do
    Agent.update(collector, fn seen ->
      [{message.data.partition, message.data.sequence, message.metadata.attempt} | seen]
    end)

    message
  end
end

{:ok, _pipeline} =
  Pipelines.RiskBroadway.start_link(stream: stream, collector: collector)

Process.sleep(1_500)

{:ok, lag} = Rheo.lag(stream, "broadway")

%{
  processed: collector |> Agent.get(&Enum.reverse/1) |> length(),
  sample: collector |> Agent.get(&Enum.reverse/1) |> Enum.take(5),
  frontiers: Map.new(lag.partitions, fn {p, info} -> {p, info.frontier} end),
  lag: lag.lag
}

Failure modes (for your handlers)

# Broadway.Message.failed(message, :reason)
#   → Rheo.nack/3  (retry until the group's max_attempts, then dead-letter)
#
# Broadway.Message.configure_ack(message, on_failure: :reject)
#   → Rheo.reject/3 (dead-letter immediately)
:ok = Broadway.stop(Pipelines.RiskBroadway)

Two knobs people confuse:

  • :max_demand on the producer — caps unsettled leases
  • Broadway :concurrency — caps parallel handler work

Choosing a surface

Surface Use when
Rheo.Consumer Simple OTP handlers, most apps
Rheo.Producer + GenStage Custom demand topology
Rheo.Producer + Flow Parallel map / partition / reduce / window
Rheo.Producer + Broadway Processors, batchers, rate limits

Rheo keeps leases, fencing, frontier, and receipts. Pipelines must not invent a second delivery identity.

Next: Backends

HexDocs: GenStage · Broadway · Flow spike · Article 14 · Article 16