Skip to content
LakeBench
ProblemsCommunityPricing
Sign inStart practicing
Back
  1. Home
  2. Interview prep
  3. Consumer lag keeps growing

Kafka · Operations & Scenarios

Consumer lag keeps growing

Hardkafka-48
consumer-lagperformance-troubleshootingpartition-skewscenario

Question

Consumer lag on a topic keeps growing. How do you investigate and fix it?

Solution

When consumer lag increases continuously, start by determining whether the lag is evenly distributed across all partitions or isolated to one or two partitions. Partition-specific lag indicates key skew or a stalled consumer thread, while uniform lag across every partition means aggregate incoming message velocity exceeds consumer processing throughput. Remediation requires eliminating slow per-record external calls, scaling out consumer instances up to the partition count, batching sink writes, and preventing rebalance storms caused by max.poll.interval.ms expirations.

Step 1: Inspect lag distribution across partitions

Query the consumer group status using administrative tools to isolate the pattern:

kafka-consumer-groups --bootstrap-server broker:9092 \
  --describe --group order-sync-service

A review of the command output reveals the underlying issue:

  • If only partition 4 has 500,000 lagging records while partitions 0 through 3 have near-zero lag, you have a hot partition key. For example, a large merchant account generating bulk orders may be hashing to a single partition.
  • If all partitions show steadily climbing lag, total consumer capacity is simply too low for the current produce rate.

Step 2: Profile per-record processing bottlenecks

Inspect consumer application logs and APM traces to check how time is spent inside the processing loop. A classic junior engineer mistake is executing synchronous network I/O per record, such as making an individual REST API call or executing a single-row SQL insert for every event. If each record takes 10 milliseconds, a single thread can process at most 100 records per second. Refactoring the consumer to buffer records in memory and execute micro-batched bulk inserts (for example, inserting 500 rows per SQL statement) can boost throughput by twenty times without adding hardware.

Step 3: Scale consumer concurrency and partitions

If existing consumer pods have saturated their CPU and memory:

  • Add consumer instances up to the total number of partitions in the topic. If the topic has 16 partitions and you only run 4 consumer pods, scaling the consumer group to 16 pods instantly quadruples reading capacity.
  • If consumer instances already equal the partition count, you cannot add more consumers to that group. You must increase the topic partition count (for instance, from 16 to 32) and deploy matching consumer instances.

Step 4: Stop rebalance storms from processing timeouts

If consumers poll large batches and spend too long processing them, they exceed max.poll.interval.ms. When this happens, the coordinator marks the consumer dead, revokes its partition assignments, and triggers a group-wide rebalance. The evicted consumer then attempts to reconnect, setting off an endless rebalance storm where no work gets done. Reduce max.poll.records from 500 down to 50 or 100 so each polled batch completes well within the timeout threshold, or raise max.poll.interval.ms to give consumers adequate processing headroom.

PreviousNext