How would you design a consumer-lag alert that pages on-call only for real problems, and what makes raw lag a poor threshold?
answer
- Lag = log-end-offset − committed offset
- Raw lag meaningless without throughput
- Time-to-drain = lag / consume rate
- Burrow: OK/WARN/ERR/STALL, no fixed threshold
- Alert on growing + sustained, and on stalled commits
basics
~20 sRaw lag (offsets behind) is misleading because acceptable lag depends on throughput. Alert on time-to-drain (lag / consume rate) or sustained-and-growing lag, not a single absolute number, and require the condition to persist before paging.
solid answer
~50 sConsumer lag = log-end-offset minus the consumer group's committed offset per partition: how many messages are unprocessed. Raw absolute lag is a poor page threshold because '1M messages behind' is fine for a high-throughput stream and catastrophic for a low one — the meaningful quantity is *time*: estimated lag-in-seconds = lag / consumption rate (this is what Burrow and the Kafka lag exporter compute). Better signals: (1) lag that is both high and *growing* over a window (consumer can't keep up), and (2) time-to-drain exceeding your freshness SLO. Burrow's evaluation goes further — it classifies a group as OK/WARN/ERR/STALL by looking at whether offsets are committing and whether lag is monotonically increasing, avoiding fixed thresholds entirely. Require the condition to be sustained (e.g. 10-15 min) so a brief rebalance or traffic spike doesn't page. Also alert on a *stalled* group — committed offset not advancing at all — which is often worse than slowly growing lag.
go deeper
Know lag = end offset − committed offset, and that growing lag means the consumer is falling behind.
Explain why absolute lag is a bad threshold and convert to time-to-drain or a growing-lag rule with a sustained window.
Design a freshness SLO, detect stalled groups vs slow groups, and choose Burrow-style trend evaluation over fixed thresholds.
Set org-wide consumer-freshness SLOs and lag-monitoring standards, and teach teams to alert on trend + commit progress, not raw counts.
## What consumer lag is Kafka stores messages in partitions as an append-only log; each message has a monotonically increasing **offset**. A **consumer group** reads partitions and periodically **commits** the offset of the last message it has processed. **Lag** for a partition = `log-end-offset (latest produced) − committed offset (last processed)` = number of messages produced but not yet processed. Total group lag is the sum across its partitions. Why it matters: lag is the primary measure of consumer *freshness*. Growing lag means downstream data is getting stale — a fraud check, a dashboard, or a cache is falling behind reality. ## Why raw lag is a poor threshold Absolute lag has no universal meaning: - A topic doing 500k msg/s can sit at 2M lag and drain it in 4 seconds — totally healthy. - A topic doing 100 msg/s at 2M lag is ~5.5 hours behind — a serious incident. So a fixed number like `lag > 100000 → page` either misses real incidents or fires constantly. The meaningful quantity is **time-to-drain = lag / consumption_rate** (lag expressed in seconds of delay). Tools like **Burrow** and **kafka-lag-exporter** compute this lag-in-seconds. ## Better alerting designs 1. **Time-based SLO.** Define a freshness SLO ('p99 of records processed within 60s of production'). Alert when estimated lag-seconds breaches it for a sustained window. 2. **Lag growing over time.** Compute the slope: is lag monotonically increasing across N consecutive checks? A rising slope means the consumer's throughput < producer's throughput — it will never catch up without intervention. 3. **Burrow's sliding-window evaluation.** Burrow keeps a window of recent (offset, lag) samples per partition and classifies the group: OK; WARN if lag is increasing; ERR if increasing and offsets aren't being committed; STALL if the committed offset isn't moving at all while lag exists. This is threshold-free and adapts to throughput. 4. **Stalled group detection.** A consumer that has stopped committing entirely (crashed, stuck in a rebalance loop, or deadlocked) may show flat lag, but it's fully broken. Alert when committed offset hasn't advanced for X minutes despite available messages. ## Avoiding false pages - Require sustained breach (10-15 min) — rebalances, deploys, and brief producer bursts cause transient lag. - Scope alerts per group + topic, and ideally per critical partition, so one hot partition is visible. - Suppress during known batch/backfill jobs that intentionally create lag. ## Edge cases - **Offset reset / new group** can show enormous lag instantly (reading from earliest) — exclude or special-case. - **Idle low-traffic topics**: lag stays near zero but a small absolute lag may still breach a tight time SLO if production trickles in. - **Rebalancing storms**: frequent rebalances cause sawtooth lag; alert on rebalance rate separately. - Measuring lag requires reading both the topic's end offset and the group's committed offset (from `__consumer_offsets`); the admin client / `kafka-consumer-groups.sh --describe` exposes both.
- A group shows flat, non-growing lag of 50k. Is that healthy?Not necessarily. Flat lag can mean a healthy steady state, but if the committed offset isn't advancing at all the group is stalled (crashed/stuck rebalance) and is fully broken — check whether commits are progressing, not just the lag number.
- How does Burrow avoid configuring per-topic lag thresholds?It keeps a sliding window of recent offset/lag samples per partition and evaluates trends — whether offsets are still committing and whether lag is monotonically increasing — classifying the group as OK/WARN/ERR/STALL. The verdict adapts to throughput instead of comparing against a fixed number.
saying these in an interview costs you the question
- Recommending a single global absolute-lag threshold for all topics
- Ignoring consumption rate / time-to-drain
- Treating flat lag as automatically healthy (misses stalled groups)
- Paging instantly on a lag spike without a sustained window (rebalances/deploys cause transient lag)