13.8 Lag monitoring
Consumer lag
- 100 s
- lag in time
- 7 min
- drains in
Healthy. 100s behind and closing at 5k/s of headroom. That headroom is what absorbs a traffic spike or a restart — a consumer running at exactly the produce rate has none.
"For Kafka consumers, the most important thing to monitor is the consumer lag. ... this is ONE OF THE CASES WHERE EXTERNAL MONITORING FAR SURPASSES WHAT IS AVAILABLE FROM THE CLIENT ITSELF."
Why the client metric is inadequate (restated):
- “It only represents a single partition, the one that has the most lag, so it does not accurately show how far behind the consumer is.”
- “It requires proper operation of the consumer, because the metric is calculated by the consumer on each fetch request. If the consumer is broken or offline, the metric is either inaccurate or not available.”
The correct architecture:
*"an EXTERNAL PROCESS that can watch BOTH the state of the partition on the broker, tracking the offset of the most recently produced message, AND the state of the consumer, tracking the last offset the consumer group has committed. This provides an OBJECTIVE view that can be updated REGARDLESS OF THE STATUS OF THE CONSUMER ITSELF.
This checking must be performed for EVERY PARTITION that the consumer group consumes. For a large consumer, like MirrorMaker, this may mean TENS OF THOUSANDS OF PARTITIONS."*
⚠️ Why the CLI approach doesn't scale
*"Monitoring lag like this, however, presents its own problems:
- "you must understand FOR EACH PARTITION what is A REASONABLE AMOUNT OF LAG. A topic that receives 100 messages an hour will need a different threshold than a topic that receives 100,000 messages per second."
- "you must be able to consume ALL of the lag metrics into a monitoring system and set alerts on them. If you have a consumer group that consumes 100,000 partitions over 1,500 topics, YOU MAY FIND THIS TO BE A DAUNTING TASK."
💡 Burrow — the recommended answer
*"an open source application, originally developed by LinkedIn, that provides consumer STATUS monitoring by gathering lag information for ALL consumer groups in a cluster and calculating A SINGLE STATUS for each group saying whether the consumer group is WORKING PROPERLY, FALLING BEHIND, or is STALLED OR STOPPED ENTIRELY.
IT DOES THIS WITHOUT REQUIRING THRESHOLDS by MONITORING THE PROGRESS that the consumer group is making on processing messages, though you can also get the message lag as an absolute number."*
"Deploying Burrow can be an easy way to provide monitoring for ALL consumers in a cluster, as well as in MULTIPLE clusters, and it can be easily integrated with your existing monitoring and alerting system."
| Threshold-based lag alerting | Status-based (Burrow) |
|---|---|
| Requires a threshold per partition. | Requires no thresholds. |
| Lag oscillates → flapping alerts. | Evaluates progress over a window. |
| 100,000 partitions = 100,000 thresholds to maintain. | One status per consumer group. |
| Breaks when the consumer is down. | Works when the consumer is down. |
"If there is NO other option, the
records-lag-maxmetric from the consumer client will provide at least a PARTIAL view. It is strongly suggested, however, that you utilize an external monitoring system like Burrow."