Command Palette

Search for a command to run...

Hectal
PHASE 14Advanced ~7 min· topic 2 of 5

Topic 14.2

Lag Observability: Records, Time, and Trend

In one line

Measure lag per partition and group in records and in time, and track its trend: is the consumer catching up, holding steady or falling behind? Alert on time lag against SLOs and on sustained growth, not on absolute record counts.

0/5 · 0%

Think of it like this

A queue at a ticket counter. "50 people waiting" means little unless you know how fast the queue moves: 5 minutes' wait is fine, 2 hours isn't, and a queue that's shrinking is much better than one that's growing.

Key ideas

  1. 01

    Record lag: log end offset − committed offset, per partition; group lag is the sum or max. Time lag: now − timestamp of the oldest unprocessed record (or record lag ÷ consume rate), which maps directly to "how stale are we".

  2. 02

    Trend: lag rate = d(lag)/dt. Negative means catching up; zero with non-zero lag means keeping up with a constant backlog; positive means falling behind. Estimate time-to-catch-up = lag ÷ (consume rate − produce rate).

  3. 03

    Partition view: one partition with growing lag while others are flat points to a hot key or one stuck consumer; group-level sums hide this.

  4. 04

    Tools: kafka-consumer-groups --describe, Burrow (evaluates lag windows and reports status like STALL or STOP), Prometheus exporters, managed consumer-lag metrics (for example MSK's MaxOffsetLag and EstimatedMaxTimeLag).

  5. 05

    SLOs: express freshness as time ("99% of events processed within 60 s") and alert when time lag threatens it; compare against retention to catch data-loss risk.

Code & diagrams

lag-queries.promqlpromql
# Record lag per group
sum by (consumergroup, topic) (kafka_consumergroup_lag)

# Is it growing? (records per second)
deriv(sum by (consumergroup) (kafka_consumergroup_lag)[10m:1m])

# Estimated time lag in seconds = lag / consume rate
sum by (consumergroup) (kafka_consumergroup_lag)
  / sum by (consumergroup) (rate(kafka_consumergroup_current_offset[5m]))

# Worst single partition, to catch hot keys hidden by sums
max by (consumergroup, topic) (kafka_consumergroup_lag)

Explain it without notes

01

Why is time lag a better alerting signal than record lag?

Practice

01

Write the PromQL for time lag and a lag-growth alert for one of your consumer groups.

Trade-offs

  • ↔

    Per-partition lag metrics are precise but high-cardinality; keep them for key groups and summarise the rest.

Done when you can

  • I monitor lag in records and time, per partition, with trend-based alerts.