Rheo Livebook Demo (v0.3.0)
Mix.install(
[
{:rheo, path: Path.join(__DIR__, "..")},
{:kino, "~> 0.14"}
],
config: [
rheo: [
start_on_application: false,
clock: Rheo.Clock.Frozen,
default_lease_ms: 5_000,
default_max_attempts: 3
]
]
)
Prerequisites
- Open this notebook from the repo (
notebooks/rheo_demo.livemd) soMix.installcan resolve{:rheo, path: ...} - Evaluate cells top to bottom
- No Docker required — this demo uses
Rheo.Backend.ETSby default
Rheo v0.3.0 is an embedded Elixir/OTP library: consumer groups over a searchable event log (Mongo durable, ETS ephemeral). Delivery is at-least-once — use event IDs for idempotency.
1. Start Rheo (ETS)
case Rheo.start_link(backend: Rheo.Backend.ETS) do
{:ok, _} -> :ok
{:error, {:already_started, _}} -> :ok
end
:ok = Rheo.ensure_indexes()
:ok = Rheo.ping()
Optional Mongo path: start Docker (
docker compose up -d), thenRheo.start_link(url: System.get_env("RHEO_MONGO_URL", "mongodb://localhost:27017/rheo_livebook"))and clear collections before continuing.
2. Create a stream and publish market events
stream = "market-events"
:ok = Rheo.create_stream(stream)
events =
for i <- 1..30 do
%{
type: "curve_update",
currency: Enum.at(["EUR", "USD", "GBP"], rem(i, 3)),
curve: Enum.at(["EUR-EURIBOR-6M", "USD-SOFR", "GBP-SONIA"], rem(i, 3)),
price: Float.round(1.0 + :rand.uniform() * 3.0, 4),
metadata: %{
correlation_id: "corr-#{rem(i, 7)}",
producer: "pricing-service-v3"
}
}
end
{:ok, written} = Rheo.append_batch(stream, events)
Kino.DataTable.new(
Enum.map(written, fn e ->
%{
sequence: e.sequence,
id: e.id,
type: e.type,
currency: e.payload["currency"],
curve: e.payload["curve"],
price: e.payload["price"]
}
end)
)
3. Independent consumer groups
risk and surveillance each consume the same immutable stream.
:ok = Rheo.create_group(stream, "risk")
:ok = Rheo.create_group(stream, "surveillance")
{:ok, risk_leases} = Rheo.fetch(stream, "risk", limit: 5, consumer_id: "risk-1")
{:ok, surv_leases} = Rheo.fetch(stream, "surveillance", limit: 5, consumer_id: "surv-1")
%{
risk_count: length(risk_leases),
surveillance_count: length(surv_leases),
same_event_ids?: Enum.map(risk_leases, & &1.event_id) == Enum.map(surv_leases, & &1.event_id)
}
Acknowledge surveillance; leave risk leases hanging (simulate a crash before ACK).
Enum.each(surv_leases, &Rheo.ack/1)
stale_risk = risk_leases
:ok
4. Lease expiry and redelivery
Advance the frozen clock past the lease TTL. A different risk worker can claim the same events; the stale lease can no longer ACK.
Rheo.Clock.Frozen.set(DateTime.utc_now())
Rheo.Clock.Frozen.advance(6_000)
{:ok, redelivered} = Rheo.fetch(stream, "risk", limit: 5, consumer_id: "risk-2")
stale_ack = Rheo.ack(hd(stale_risk))
fresh_ack = Rheo.ack(hd(redelivered))
Enum.each(tl(redelivered), &Rheo.ack/1)
%{
redelivered: length(redelivered),
stale_ack: stale_ack,
fresh_ack: fresh_ack,
attempts: Enum.map(redelivered, & &1.attempt)
}
5. Search history after consumption
Events stay in Mongo after ACK — that is Rheo's differentiator.
{:ok, eur} =
Rheo.query(stream,
type: "curve_update",
currency: "EUR",
curve: "EUR-EURIBOR-6M",
limit: 10
)
{:ok, by_corr} = Rheo.query(stream, correlation_id: "corr-1", limit: 10)
Kino.Layout.grid(
[
Kino.DataTable.new(
Enum.map(eur, fn e ->
%{sequence: e.sequence, price: e.payload["price"], id: e.id}
end)
),
Kino.DataTable.new(
Enum.map(by_corr, fn e ->
%{
sequence: e.sequence,
currency: e.payload["currency"],
correlation_id: e.metadata["correlation_id"]
}
end)
)
],
columns: 2
)
6. Idiomatic Rheo.Consumer
defmodule Livebook.RiskConsumer do
use Rheo.Consumer, stream: "market-events", group: "risk"
@impl true
def setup(opts) do
{:ok, %{seen: [], agent: Keyword.fetch!(opts, :agent)}}
end
@impl true
def handle_event(event, %{agent: agent} = state) do
Agent.update(agent, fn ids -> [event.id | ids] end)
{:ack, state}
end
end
{:ok, agent} = Agent.start_link(fn -> [] end)
{:ok, _pid} =
Livebook.RiskConsumer.start_link(
stream: stream,
group: "risk",
agent: agent,
name: Livebook.RiskConsumer,
max_demand: 10,
poll_ms: 100
)
Process.sleep(1_500)
seen = Agent.get(agent, &Enum.reverse/1)
%{processed: length(seen), sample_ids: Enum.take(seen, 5)}
Takeaways
- Immutable log — consumption never deletes events
- Independent groups — risk and surveillance progress separately
- Leases + fencing — stale ACKs fail; expired work is redelivered (at-least-once)
- Search — investigate history with
Rheo.query/2after processing - OTP — supervise
{Rheo, ...}andRheo.Consumerin your app
CLI alternative from the repo root: mix rheo.demo