01 / 05

A consumer's lag is increasing. What are the first things you would check?

Difficulty: 3/10
Production incidents, Consumer lag, Rebalancing

Diagnosing Increasing Consumer Lag: A Disciplined Approach

When consumer lag is increasing, the disciplined approach is to compare the produce rate to the consume rate, and then work backwards to find where the bottleneck is. Lag grows when produce rate exceeds consume rate, so the first question is: is the producer sending more, or is the consumer processing less? Check the produce rate per partition and per topic. If it has spiked, the consumer may be fine but overwhelmed by a burst. If it is steady, the consumer is the problem. Then check the consume rate: is the consumer processing fewer records per second than before? If so, the cause could be slower per-record processing, a downstream dependency that has slowed down, an increase in errors and retries, or a consumer that is blocked or paused. The second thing to check is the number of active consumers and the partition assignment. If a consumer has crashed or been removed from the group, the remaining consumers have more partitions each and may not keep up. If a rebalance is in progress, processing is paused, and lag grows during the rebalance. The third thing is errors: check the consumer's error rate for exceptions, deserialization failures, or downstream timeouts. A spike in errors often indicates a downstream problem or a poison message.

The mechanism for diagnosis is to use the metrics and tools that Kafka provides. The kafka-consumer-groups.sh tool shows per-partition lag, which tells you whether the lag is uniform or concentrated in a few partitions. If it is concentrated, the cause is likely skew or a slow partition. The consumer's JMX metrics show records-consumed-rate, records-lag-max, and fetch-latency-avg. The broker's metrics show produce rate, fetch rate, and request latency. The downstream service's metrics show response time and error rate. The order of checks should be: produce rate, consume rate, partition assignment, error rate, downstream latency, and resource usage (CPU, memory, disk, network) of the consumer instances. A common mistake is to jump to scaling the consumer group without checking whether the partition count allows more consumers. If the topic has 10 partitions and the group already has 10 consumers, adding more will not help. Another mistake is to ignore the possibility of a hot partition; if one partition has most of the lag, the issue is key skew, not consumer capacity. The trade-off is between speed and thoroughness. In an incident, you want to find the cause quickly, but a rushed diagnosis can lead to the wrong fix. A checklist approach ensures you cover the common causes in order.

There are several common causes of increasing lag that are easy to overlook. First, a consumer that is doing synchronous calls to a downstream service for each record will be slow; batching or async calls can help. Second, a consumer that is committing offsets too frequently will spend time on commits instead of processing; committing per batch is more efficient. Third, a consumer that is using a deserializer that is slow (e.g., Avro with a cold schema cache) will have higher per-record cost. Fourth, a consumer that is in a rebalance loop will spend most of its time rebalancing instead of processing. Fifth, a consumer that is blocked on a full send buffer or a slow disk (if it writes to disk) will lag. Sixth, a consumer that is using a small max.poll.records will poll frequently and process small batches, which can be inefficient. The trade-off is between latency and throughput. A consumer optimized for low latency will poll frequently and process small batches; a consumer optimized for throughput will poll less often and process larger batches. The right choice depends on the use case. Version note: Kafka's consumer group protocol has improved with cooperative rebalancing (KIP-429, Kafka 2.4+), which reduces the pause during rebalances, but rebalances still cause lag. Also, the static membership feature (group.instance.id) can reduce rebalances for stable consumers; if you are not using it, consider it.

javascript
  1. 1

    Compare produce rate to consume rate; lag grows when produce exceeds consume.

  2. 2

    Check per-partition lag to see if it is uniform or concentrated (skew).

  3. 3

    Check the number of active consumers and partition assignment; a crashed consumer or rebalance causes lag.

  4. 4

    Check error rate and downstream latency; errors and slow downstream calls reduce consume rate.

  5. 5

    Do not scale consumers beyond the partition count; it does not help.

  6. 6

    Common causes: synchronous downstream calls, frequent commits, rebalance loop, small batches, slow deserializer.

  7. 7

    Use static membership (group.instance.id) to reduce rebalances for stable consumers.

Share

Share via WhatsApp, X, Facebook, LinkedIn or copy link. Open Graph preview enabled.