Rheo is an Elixir/OTP library: supervise a Rheo instance, append events to a stream, and consume them with a durable consumer group. Delivery is at-least-once — use event IDs for idempotency.
Install
def deps do
[{:rheo, "~> 1.0"}]
endPick a backend when you start Rheo:
| Backend | When |
|---|---|
Rheo.Backend.ETS | Tests, Livebook, ephemeral apps |
{Rheo.Backend.Mongo, url: ...} | Durable MongoDB |
{Rheo.Backend.Ecto, repo: MyApp.Repo} | Durable PostgreSQL or SQLite |
Minimal ETS example
children = [
{Rheo, name: MyRheo, backend: Rheo.Backend.ETS}
]
Supervisor.start_link(children, strategy: :one_for_one)
:ok = Rheo.create_stream("orders", rheo: MyRheo)
:ok = Rheo.create_group("orders", "fulfillment", rheo: MyRheo)
{:ok, _event} =
Rheo.append("orders", %{type: "order_created", key: "cust-1"}, rheo: MyRheo)
{:ok, [lease]} = Rheo.fetch("orders", "fulfillment", limit: 1, rheo: MyRheo)
:ok = Rheo.ack(lease, rheo: MyRheo)Idiomatic consumer
defmodule MyApp.FulfillmentConsumer do
use Rheo.Consumer, stream: "orders", group: "fulfillment"
@impl true
def handle_event(event, _context) do
:ok = MyApp.Fulfillment.process(event)
:ack
end
end
children = [
{Rheo, name: MyRheo, backend: Rheo.Backend.ETS},
{MyApp.FulfillmentConsumer, rheo: MyRheo, concurrency: 4, max_demand: 50}
]Next steps
- Configuration — leases, demand, clocks, auto-start
- Consumer Groups — competing vs independent groups
- ETS / Mongo / Using Ecto
- Broadway / GenStage
- Livebook demos — Quickstart, Concepts, Pipelines, Backends