From v0.5, a stream may declare multiple partitions. Ordering and sequences are per partition; there is no global order across partitions.
Create a partitioned stream
:ok = Rheo.create_stream("market-events", partition_count: 4)Append with a routing key (or explicit :partition):
{:ok, event} =
Rheo.append("market-events", %{type: "curve_update", key: "EUR-1", price: 1.0})
event.partition
event.sequenceKey routing uses :erlang.phash2/2. Same key → same partition.
Contiguous frontier
Group progress is the highest contiguous terminal sequence per partition. ACK holes do not advance the frontier:
ACK 1001, inflight 1002, ACK 1003 → frontier stays at 1001{:ok, lag} = Rheo.lag("market-events", "risk")
lag.partitions[0].frontier
lag.partitions[0].high_water
lag.lagFetch by partition
Rheo.fetch(stream, group, partition: 0, limit: 10)
# Consumer / Producer
partitions: [0, 1]
# or partitions: :allDesign: ADR 016, tutorial ACKs are not a cursor, and the partitions diagram in Architecture.