I wrote Rheo because I kept needing two things from the same events: deliver the next one to a group of workers, and look up what happened last Tuesday. Most setups solve that with a broker plus a database. I wanted the consumer-group machinery in OTP, on a store I was already running.

Rheo is an Elixir library. You put it in your supervision tree. It gives you leases, competing consumers, ACK fencing, partitions, lag, replay, and query over MongoDB, PostgreSQL, SQLite, Redis Streams, Mnesia, or ETS.

Library: hex.pm/packages/rheo · Docs: rheo.hexdocs.pm · Source: github.com/thanos/rheo

Delivery is at-least-once. Order is per partition. Consuming an event does not delete it.

Rheo example control — append events, inspect groups, watch lag

The example app: append to the log, inspect groups, watch lag.

Parts 1–8 are the idea and the mechanics. Parts 9–17 are why the API moved: the 0.2 break, ETS as a second witness, Redis forcing a rewrite of what the callbacks mean, and what I froze before 1.0.

If you want to try it, mix rheo.demo and the Livebook demos are the hands-on version. HexDocs is the contract. These articles are the argument.

I wrote them so I could remember why I made the choices I did, and so you can decide whether the library is for you.

Articles in this series

  1. Rheo: Why Put Consumer Groups in Front of a Database? — I wanted to deliver events to a group of workers and still be able to look those events up later, without running two systems.
  2. Rheo: What Is a Consumer Group? — A stream is the log. A group is independent progress on that log. Workers in a group compete. Groups do not.
  3. Rheo: Why ACK Is Harder Than It Looks — Acknowledging a message looks like one line. The failure cases are the part I actually had to design.
  4. Rheo: Building Rheo as an Elixir/OTP Library — Rheo is not a server you run. It is a child spec you supervise. The database keeps the truth.
  5. Rheo: Demand, Backpressure, and Database Consumers — If you fetch without a bound, a crash turns into a redelivery storm. I use max_demand for that.
  6. Rheo: MongoDB as a Searchable Event Log — The first backend was Mongo because I wanted to query the history, not just pop a queue.
  7. Rheo: Killing Consumers on Purpose — If killing the worker loses the message, you built a demo. I like a rude one.
  8. Rheo: Searching the Stream — After ACK I still want to ask what happened. The log is still there.