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:
- A plain GenStage consumer
- Flow map / reduce / window (with settlement rules)
- 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}— eitherRheo.ConsumerorRheo.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:
- Settle durably (
Rheo.ack/1orProducer.ack/3) - Free the producer's inflight slot (
confirm/2is included inProducer.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_demandon 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