Rheo.Lag (rheo v1.0.0)

Copy Markdown View Source

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 3

Fields

FieldTypeMeaning
streamString.t()Stream name
groupString.t()Consumer group name
partitions%{partition => partition_lag()}Per-partition breakdown
lagnon_neg_integer()Sum of per-partition lags

partition_lag map

KeyMeaning
:frontierContiguous committed sequence for the group on that partition
:high_watermarkHighest sequence written on that partition
:lagmax(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.

t()

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

partition_lag()

@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.

t()

@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

build(stream, group, partitions)

@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
3

Returns

A %Rheo.Lag{}. Does not raise when partitions is a map.

from_maps(stream, group, frontiers, high_watermarks, partition_list)

@spec from_maps(String.t(), String.t(), map(), map(), [non_neg_integer()]) :: t()

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.