Per-partition and aggregate consumer lag.
Lag for a partition is max(high_watermark - frontier, 0). Aggregate lag
is the sum of per-partition lags. Hosts normally call Rheo.lag/3;
build/3 and from_maps/5 are for backends and tests assembling the struct.
partition 0: HW=10 frontier=7 → lag 3
partition 1: HW=5 frontier=5 → lag 0
aggregate lag 3Fields
| Field | Type | Meaning |
|---|---|---|
stream | String.t() | Stream name |
group | String.t() | Consumer group name |
partitions | %{partition => partition_lag()} | Per-partition breakdown |
lag | non_neg_integer() | Sum of per-partition lags |
partition_lag map
| Key | Meaning |
|---|---|
:frontier | Contiguous committed sequence for the group on that partition |
:high_watermark | Highest sequence written on that partition |
:lag | max(high_watermark - frontier, 0) |
Example
iex> lag = Rheo.Lag.build("market-events", "risk", %{
...> 0 => %{frontier: 7, high_watermark: 10, lag: 3},
...> 1 => %{frontier: 5, high_watermark: 5, lag: 0}
...> })
iex> {lag.lag, lag.partitions[0].lag}
{3, 3}
Summary
Types
Lag numbers for one partition: frontier, high-watermark, and their difference.
Aggregate lag for a {stream, group} with a per-partition map.
Functions
Builds a lag struct from per-partition frontier/HW maps.
Computes lag entries for each partition given frontier and high-watermark maps.
Types
@type partition_lag() :: %{ frontier: non_neg_integer(), high_watermark: non_neg_integer(), lag: non_neg_integer() }
Lag numbers for one partition: frontier, high-watermark, and their difference.
@type t() :: %Rheo.Lag{ group: String.t(), lag: non_neg_integer(), partitions: %{optional(non_neg_integer()) => partition_lag()}, stream: String.t() }
Aggregate lag for a {stream, group} with a per-partition map.
See the module documentation for field meanings.
Functions
@spec build(String.t(), String.t(), %{optional(non_neg_integer()) => partition_lag()}) :: t()
Builds a lag struct from per-partition frontier/HW maps.
Sums each entry's :lag into the aggregate. Does not recompute lag from
frontier/HW — pass already-calculated values (see from_maps/5).
Examples
iex> Rheo.Lag.build("s", "g", %{0 => %{frontier: 1, high_watermark: 4, lag: 3}}).lag
3Returns
A %Rheo.Lag{}. Does not raise when partitions is a map.
Computes lag entries for each partition given frontier and high-watermark maps.
Missing keys default to 0 via Rheo.Partition.map_get/3.
Examples
iex> lag = Rheo.Lag.from_maps("s", "g", %{"0" => 2}, %{"0" => 5}, [0])
iex> lag.partitions[0]
%{frontier: 2, high_watermark: 5, lag: 3}Returns
A %Rheo.Lag{} covering every id in partition_list.