Kafka Consumer Lag Growing but CPU Idle — Fix It
Lag climbs while brokers sit idle because your consumer is slow.
20+ years shipping production backend systems. Written from production experience, not tutorials.
- ✓A running Kafka consumer group you can describe
- ✓Basic consumer configs (poll records, intervals)
- ✓Access to consumer logs and JMX metrics
- Lag with idle brokers means your consumer is the bottleneck — brokers can't push faster than you poll and process
- Cut max.poll.records to 50-200 so each batch finishes well inside max.poll.interval.ms instead of timing out
- Profile the handler first: one synchronous DB call per record turns 500 records into a batch that never finishes
- Confirm with kafka-consumer-groups.sh --describe plus GC logs and JMX fetch metrics before scaling anything
Picture a warehouse conveyor belt feeding packing tables. The belt (Kafka broker) runs smooth and quiet, but boxes pile up because one packer opens every box, phones the supplier, waits on hold, then packs it. The belt isn't broken — the packer is slow. Speeding up the belt changes nothing. You fix it by giving the packer smaller stacks of boxes, a faster phone routine, and extra hands — that's exactly what tuning max.poll.records, handlers, and consumer count does.
It's 10 AM and the lag dashboard is a ski slope pointing the wrong way: 2.4 million messages behind and climbing 40,000 every minute. You pull up the broker dashboards expecting fire and find nothing — CPU at 15%, disks calm, network flat. Somebody says the cluster is fine and closes the tab. They're half right, and that half will cost you the afternoon.
Idle brokers with growing lag is Kafka's classic misdirection. Brokers only store and serve; they can't force your consumer to poll faster. When lag grows and brokers yawn, the bottleneck lives in your consumer process: a handler that takes 800 ms per record, a batch sized for a benchmark you never ran, a garbage collector freezing the world for 30 seconds, or a rebalance loop replaying the same records forever.
This guide walks the exact diagnosis order that separates those four causes in minutes. You'll learn to read per-partition lag instead of totals, size max.poll.records from measured timings, spot GC stalls in logs, and confirm rebalance storms with two greps. By the end, idle brokers won't fool you again — you'll know precisely which consumer-side knob to turn.
Why Lag Grows While Brokers Look Idle
Consumer lag is the gap between what producers appended and what your group committed: LAG = log-end-offset minus committed offset, per partition. Brokers track both numbers, but neither number measures broker effort. Serving a fetch is cheap — a sequential disk read the page cache usually handles. Processing the fetched records is where real work happens, and that work lives entirely in your consumer process.
That's why idle brokers prove nothing about consumer health. A broker at 15% CPU happily serves 50 MB/s of fetches while your consumer chews one record per second through a blocking HTTP call. The end offset races ahead, commits crawl, and the subtraction grows. Every minute of slow processing adds 60 seconds of catch-up debt, and the debt compounds because new records keep arriving while you grind through old ones.
The mental shift that fixes this class of incident: treat brokers and consumers as separate systems with separate dashboards. Broker CPU, disk, and request latency tell you about the cluster. Per-partition lag, batch processing time, and commit rate tell you about the consumer. When the first set is green and the second is red, stop touching the cluster — you'll find the answer in poll-loop timings, handler profiles, and GC logs.
max.poll.records Too High: Big Batches, Slow Commits
max.poll.records caps how many records one poll() call returns — default 500. max.poll.interval.ms is the deadline to call poll() again before the coordinator declares you dead — default 300,000 ms (5 minutes). These two knobs form a contract: everything you fetched must be processed and the next poll() issued before the interval expires. Break the contract and your partitions get revoked mid-batch.
Do the arithmetic openly. If your handler needs 1.2 seconds per record (a realistic number with one downstream call), 500 records need 600 seconds — double the default interval. The consumer looks busy the whole time, brokers look idle, and the group rebalances every 5 minutes forever. Cut max.poll.records to 100 and the same work needs 120 seconds, comfortably inside the limit with room for GC wobble.
Size from measurements, not vibes. Time your slowest realistic batch on production-shaped data — big records, cold caches, slow downstream — then pick max.poll.records so the worst case lands under half the interval. That 2x headroom absorbs traffic spikes and minor GC pauses without a revoke. Re-measure after every handler change, because one new API call per record can silently double batch time and restart the whole cycle.
Slow Consumer Processing: Finding the Hot Path
Slow processing is the most common root cause and the hardest to see from Kafka dashboards, because the consumer never errors — it just takes forever. The usual villain is blocking I/O inside the record loop: a pricing API call, a Postgres upsert, an S3 read, each innocent alone and lethal at 500x. One 300 ms call per record caps you at ~3 records/sec; a 1,000-record batch then needs over 5 minutes and the interval kills you first.
Profile before optimizing. Log batch start, batch end, and record count on every poll cycle — three lines that immediately reveal per-record cost. If per-record time dominates, attack it structurally: batch the writes (200-row inserts instead of 200 single upserts), cache the hot reads (a 5-minute TTL on pricing data hit 94% in one real incident), or fan CPU-bound work to a bounded thread pool while keeping offset commits strictly ordered after the full batch.
Keep one correctness rule sacred: commit offsets only after the entire batch is durably handled. Parallelizing handlers tempts people into committing early, which trades duplicates (annoying, recoverable) for data loss (silent, permanent). Slow-but-correct beats fast-and-lossy every time — then make correct fast with batching and caching.
GC Pauses That Stall the Poll Loop
The JVM can stall your consumer without executing a single line of your code. A full garbage collection freezes all threads — including the background heartbeat thread that tells the coordinator you're alive. Miss enough heartbeats past session.timeout.ms and you're evicted: partitions revoke, lag jumps, and the next consumer replays everything you processed but never committed.
These stalls are invisible in application logs because nothing runs during the pause — the log just resumes 30 seconds later as if nothing happened. The evidence lives in GC logs, which most teams don't ship. A consumer on a default 1 GB heap doing heavy JSON parsing can easily hit 20-60 second full GCs under load, and each one is a guaranteed rebalance.
Treat GC observability as mandatory consumer infrastructure. Run every consumer JVM with unified GC logging to a shipped file, set -Xms equal to -Xmx so the heap never resizes under load, and prefer G1GC or ZGC for pause-sensitive poll loops. Alert on any single pause over 5 seconds — that's half a session timeout gone in one freeze. When heartbeat failures correlate with GC pauses in the same minute, you've found your culprit and no handler rewrite will help until the heap is fixed.
Rebalance Storms That Reset Your Progress
Rebalances turn slow consumers into stuck consumers. Every revoke discards processed-but-uncommitted work, and the next owner replays it from the last commit — paying the slow-handler cost all over again. With the default eager assignor, one member leaving pauses the entire group; with rolling deploys on 12 pods, the group can spend more time rebalancing than processing, and lag climbs in the famous sawtooth: drains a little, jumps back, drains, jumps.
Two settings end most storms. Static membership (group.instance.id set to a stable pod-unique value) lets a restarted member reclaim its partitions without a full reshuffle — the coordinator recognizes it instead of treating it as a stranger. The cooperative sticky assignor migrates only the partitions that must move instead of revoking everything, so the group keeps processing through membership changes.
Verify calm the same way you diagnosed the storm: count 'Revoked partitions' lines per hour (steady groups should show near zero) and watch per-partition lag drain monotonically after a deploy. If a rolling restart still causes a group-wide stop, one consumer is missing its group.instance.id — grep the configs, because a single dynamic member can force eager behavior for everyone.
Scaling Out: Partitions, Threads, and Fetch Tuning
Scaling out is the last step, not the first — but when the handler is tuned and batches fit, it's the right one. The hard ceiling is partition count: a group can't assign one partition to two consumers, so the 7th consumer on a 6-partition topic idles forever while you pay for its pod. Check kafka-topics --describe first; if consumers already equal partitions, add partitions (accepting the key-ordering change) before adding pods.
Use JMX fetch metrics to confirm headroom exists before scaling. records-lag-max rising with a healthy fetch-rate means records arrive fine and processing is the limit — more consumers help. Near-zero fetch-rate means records aren't arriving (throttling, network, fetch config), and extra consumers will idle just like the current ones. Scale the bottleneck you measured, not the one you assumed.
Roll out gradually and watch the drain rate, not the absolute lag. Double consumers, confirm records/sec roughly doubles and per-partition lag falls evenly. If one partition stays hot while others drain, you've got key skew — no consumer count fixes a single key holding 80% of traffic. Fix the partition key or split the hot entity before burning more pods. Start with one extra consumer per hot partition, confirm the drain rate climbs linearly, and stop when per-partition lag falls evenly — linear gains mean you found the true ceiling.
The 800 ms Lookup That Built a 2.4M Lag Behind Idle Brokers
- Time one batch before changing any infrastructure — 1,000 records at 800 ms each can never fit a 5-minute poll interval, and no broker size fixes arithmetic.
- Per-record external calls are the usual killer; a cache with 94% hit rate did more than doubling broker hardware ever could.
- Set static membership on every consumer from day one — without group.instance.id, each deploy replays hours of already-processed records.
poll() and around the commit, then compute records * seconds-per-record. If 500 records at 1.2s each equals 600s against max.poll.interval.ms=300000 (300s), the batch can never finish. Cut max.poll.records to 50 and re-measure until the worst-case batch lands under 150s — half the interval.| File | Command / Code | Purpose |
|---|---|---|
| consumer.properties | max.poll.records=100 | max.poll.records Too High |
| consumer_poll_loop.py | from kafka import KafkaConsumer | Slow Consumer Processing |
| gc_check.sh | export KAFKA_HEAP_OPTS="-Xms2g -Xmx2g" | GC Pauses That Stall the Poll Loop |
| membership_fix.sh | grep -cE 'Revoked partitions|Attempt to heartbeat failed' /var/log/app/consumer.... | Rebalance Storms That Reset Your Progress |
| scale_check.sh | kafka-topics.sh --bootstrap-server $BROKERS --describe --topic recs | Scaling Out |
Key takeaways
Common mistakes to avoid
5 patternsLeaving max.poll.records at the 500 default with slow per-record work
Doing synchronous HTTP or DB calls per record inside the poll loop
Running consumers on default 1 GB heaps with CMS or no GC logging
Redeploying consumers without static membership
Adding consumer instances when partitions are the real ceiling
Interview Questions on This Topic
Lag climbs while broker CPU sits at 15%. What does that tell you?
Frequently Asked Questions
20+ years shipping production backend systems. Written from production experience, not tutorials.
That's Kafka. Mark it forged?
5 min read · try the examples if you haven't