skip to content

You're on call for a CQRS system where a projector consumes a Kafka topic to update a read-model database. Support reports that some users see stale search results for several minutes after editing a listing. What would you measure to confirm and bound the lag window, and what are two concrete mitigations if the backlog keeps growing under load?

level: seniorimportance: should knowfreq 55%

answer

  1. consumer/offset lag vs apply/processing lag
  2. canary probe write-then-poll
  3. partition and consumer scale-out
  4. batch writes to lower per-event cost
  5. flat lag under continued writes = poison pill, not throughput

basics

~20 s

Measure the gap between when an event was written and when the projector actually applied it (processing lag), plus how many unprocessed events are queued up (offset lag). If the backlog keeps growing, either process events faster with more parallel consumers, or make each projector write cheaper, for example by batching, so it can keep pace.

solid answer

~50 s

Instrument end-to-end lag directly: stamp a timestamp when the write-side event is produced and again when the projector applies it, then track apply_time minus produce_time as an alertable p99 latency metric. Separately watch Kafka consumer-group offset lag, how far behind the topic's high-water mark the consumer sits, as a leading indicator of backlog even before it shows up as end-to-end latency. If the backlog keeps growing under load, meaning production rate exceeds consumption rate, two concrete fixes: (1) scale out consumption by increasing topic partitions and running more projector instances in the same consumer group so work parallelizes; (2) reduce per-event processing cost, such as batching writes to the read store instead of one write per event, or moving a slow enrichment step to an async secondary pass. A synthetic canary write, timed until it's visible via the query API, gives an outside-in measurement of what users actually experience.

go deeper

for a junior

Not expected to design monitoring; should understand that 'lag' is something you can and should measure, not just guess about.

for a middle

Should know Kafka-style consumer lag is a standard metric and be able to read a lag dashboard, even without designing the alerting strategy.

for a senior

Should instrument both offset lag and end-to-end apply lag, diagnose a growing backlog to root cause (throughput imbalance vs poison pill vs downstream bottleneck), and apply an appropriate mitigation.

for a principal

Should set the org's SLOs for read-model freshness, choose default instrumentation and alerting patterns for all projector pipelines, and decide when growing backlog risk justifies architectural change over tuning.

## The two signals worth measuring There are two distinct lag signals worth measuring, and conflating them is a common mistake. | Signal | What it is | What it tells you | |---|---|---| | Consumer (or offset) lag | a broker-side, queue-depth signal: how many messages are sitting unread in the topic, measured as the gap between the consumer group's committed offset and the topic's current high-water mark | it's cheap to expose (Kafka's own tooling, or Prometheus exporters like Burrow, surface it directly) and it's a leading indicator — it can start climbing before any individual message is slow to process | | Apply lag, sometimes called end-to-end or processing lag | the actual time delta between when the event was produced and when the projector finished applying it to the read store | this is the metric closest to what users actually experience as staleness, but by itself it can look fine right up until a backlog suddenly overflows it | A synthetic **canary probe** complements both: periodically write a marker record with a known timestamp and poll the query API until it becomes visible, giving a real, outside-in measurement independent of which real records happen to get read back soon after being written. ## Why the distinction matters This distinction matters because CQRS's entire value proposition rests on the promise 'we accept some delay in exchange for read-side scalability', and without measuring both signals you can't actually verify that promise is being kept versus quietly broken. - A growing backlog under sustained load specifically indicates a structural throughput imbalance — consumption rate is now below production rate — which won't self-correct; it will keep getting worse until either load drops or the system is scaled or optimized. - Distinguishing a stock metric (backlog depth) from a rate metric (lag time) is exactly the difference between noticing the leak early versus noticing it only once the bucket overflows. ## The two mitigations The two mitigations trade off differently. 1. **Scaling out consumption** — adding partitions and running more consumer instances within the group — adds real infrastructure cost and requires the topic to already be partitioned with enough headroom; repartitioning a live topic is disruptive, and because Kafka's ordering guarantee only holds within a single partition, changing the partition count or the partitioning key can change which messages are guaranteed to arrive in order relative to each other. 2. **Batching writes to the read store** reduces per-event overhead and can meaningfully raise sustainable throughput, but it introduces a minimum latency floor — you're now waiting to accumulate a batch before writing — and complicates error handling, since one bad record inside a batch needs its own retry or dead-letter path rather than just failing the whole batch. ## Failure modes mistaken for ordinary lag Several failure modes commonly get mistaken for ordinary lag. - **A 'poison pill' event** — one the projector can't process due to a bad schema or an unexpected null field — can stall an entire partition indefinitely if there's no dead-letter handling, which looks superficially like severe lag but is actually a full stop that will never recover on its own; a flat, unchanging offset lag despite continued upstream writes is the telltale signature of this, versus a steadily climbing lag which indicates genuine throughput imbalance. - **Consumer-group rebalances** triggered by scaling changes cause brief lag spikes as partitions get reassigned, which is normal and self-resolving and shouldn't trigger a page. - **A slow downstream dependency inside the projector itself** — a synchronous enrichment call to another service, for instance — can silently throttle throughput well before generic CPU or memory metrics show anything unusual, which is why timing each internal step of the projector, not just the overall consumer loop, is worth doing. - **Alerting that isn't tuned** to distinguish 'elevated but recovering during normal peak load' from 'growing without bound' produces alert fatigue that eventually gets the whole lag alarm ignored. ## The standard tooling This class of problem is common enough to have standard tooling built around it: Kafka's own consumer-lag metric (`records-lag-max`) is the de facto standard instrument for offset lag, and pipelines built on Debezium plus Kafka Connect — a common way to stream row-level changes from a write-side relational database into a search index like Elasticsearch for a CQRS read model — routinely expose 'connector lag' as a first-class operational SLO precisely because search-index freshness is exactly the kind of user-facing symptom described here.

  • What's the difference between consumer/offset lag and end-to-end/apply lag, and why track both?
    Offset lag is a queue-depth signal - how many messages are unread - cheap to measure straight from the broker and a leading indicator. Apply lag is the actual user-felt latency from write to read-model visibility, including processing time. A growing offset lag with stable apply lag can mean processing has just started falling behind; watching only apply lag can miss a backlog that's already building before it manifests as slow reads.
  • Why might adding more partitions not immediately fix a growing backlog?
    More partitions only help if there are idle consumer instances able to pick up the new ones - consumer count has to scale too. Repartitioning also changes the key-to-partition mapping going forward, which can affect ordering guarantees for keys that matter. And a slow downstream write target, like a database at its I/O ceiling, can remain the bottleneck regardless of how parallel the consumers are.
  • How does a canary/synthetic probe differ from just averaging real user-facing lag metrics?
    A canary writes a known marker at a known time and polls for its visibility, giving a clean, controlled, always-on measurement independent of real traffic volume or which records users happen to re-read soon after writing them. Real-user aggregate metrics can be noisy and biased, since many read-model rows are never queried again soon after a write, making them harder to alert on with a clear threshold.

Like tracking a mail sorting facility with two gauges: 'how many bags are sitting unsorted on the floor' (backlog/offset lag) and 'how long from drop-off to delivery' (end-to-end lag) - a growing pile on the floor is the leading warning sign that the delivery-time gauge is about to start climbing too.

saying these in an interview costs you the question

  • only mentions CPU/memory metrics with no lag-specific instrumentation
  • conflates a stalled poison-pill consumer with generic 'slow' lag
  • proposes only 'add more servers' with no distinction between broker-side backlog and processing throughput
  • no mention of measuring end-to-end, user-felt latency at all, only infra-internal metrics

context