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.