Learn Labs
4. Kafka Consumers: Reading Data from Kafka

4.13 What actually breaks in production — Ch. 4 consolidated

Production failure catalog
0 rows
#SymptomRoot causeFix
1Consumer group falls behind and never catches upSingle consumer; high-latency per-record work (DB write)Add consumers up to the partition count; beyond that, add partitions (mind Ch. 3 §9.4)
2Added consumers, throughput unchangedMore consumers than partitions → the extras are idlePartition count is the ceiling
3A second application "steals" messages from the firstBoth apps used the same group.idOne group.id per application
4Rebalance storm — group never stabilizesProcessing per batch exceeds max.poll.interval.ms → evicted → rebalance → new owner also evictedLower max.poll.records; raise max.poll.interval.ms; move slow work off the poll thread
5Entire group pauses for seconds on every deployEager rebalance = "stop the world"CooperativeStickyAssignor (mind the <2.3 upgrade path)
6Consumers assigned very uneven partition countsRangeAssignor is the default and assigns per topic independentlyRoundRobinAssignor / StickyAssignor / CooperativeStickyAssignor
7session.timeout.ms dead air after every restartclose() not called → coordinator must wait for session timeoutShutdown hook → wakeup() → close()
8Local cache/state rebuilt on every restart, taking minutesDynamic membership → new member ID → new partitionsgroup.instance.id (static membership)
9After enabling static membership, partitions sit unconsumed after a crashStatic members don't leave proactively; detection waits for session.timeout.msTune session.timeout.ms: high enough to survive restarts, low enough to reassign on real failure
10group.instance.id already exists errorTwo consumers with the same group.instance.idUnique per instance
11Duplicate processing after every crashAutocommit — last commit up to auto.commit.interval.ms (default 5 s) oldManual commits after processing; duplicates can be reduced but never eliminated with autocommit
12Silent data loss when processing throws mid-batchAutocommit + catch+continue → next poll() commits the whole batch including unprocessed records"critical to always process all events returned by poll() before calling poll() again"; use manual commits
13Records committed but never processed (loss)commitSync() called before finishing the batchCommit after processing
14Committed offset went backwards; duplicates spikedcommitAsync() retried in a callback, landing an older offset after a newer oneDon't retry naively; use the monotonic sequence-number pattern
15Last batch reprocessed after every clean shutdownOnly commitAsync() used; the final commit failed with no "next commit" to retry itcommitAsync() in the loop + commitSync() on exit
16Batch reprocessed after every rebalanceNo offsets committed before losing partitionscommitSync(currentOffsets) in onPartitionsRevoked()
17Off-by-one: one record duplicated (or skipped) per commitCommitted record.offset() instead of record.offset()+1The committed offset is the next offset to read
18A consumer down for a long weekend silently skips everythingCommitted offset aged out of retention + auto.offset.reset=latest (the default)Set earliest (accept reprocessing) or none (fail loudly); retention ≥ max outage (Ch. 2)
19A monthly batch consumer group "forgets" everythingGroup empty longer than offsets.retention.minutes (default 7 days) → offsets deleted → behaves as brand newRaise the broker setting, or keep a member alive, or manage offsets externally
20Broker/network overloaded by metadata, not dataRegex subscribe() on a cluster with ~30,000+ partitions — filtering is client-side, full topic list fetched at intervalsExplicit topic lists; "bandwidth used by topic metadata can be larger than the bandwidth used to send data"
21Regex subscription fails with authorization errorsRegex requires a full describe grant on the entire clusterExplicit lists, or grant cluster-wide describe (usually unacceptable)
22Migrating poll(0) to poll(Duration.ofMillis(0)) broke metadata fetchingpoll(long) blocked for metadata regardless of timeout; poll(Duration) does notMove the logic into onPartitionsAssigned()
23ConcurrentModificationException / bizarre consumer behaviorSharing one consumer across threads, or two same-group consumers in one threadOne consumer per thread; ExecutorService or a consumer→queue→workers pattern
24Consumer OOM after a rebalance shrinks the groupmax.partition.fetch.bytes × (now-larger) partition countUse fetch.max.bytes — an absolute cap
25High consumer CPU on an idle topicTight poll loop with a short timeout and fetch.min.bytes=1Raise fetch.min.bytes and/or the poll timeout
26Latency spikes to ~500 ms on a low-traffic topicfetch.min.bytes raised, fetch.max.wait.ms left at the 500 ms defaultLower fetch.max.wait.ms to your SLA
27Lowering request.timeout.ms made an incident worseDisconnect/reconnect churn against an already-overloaded broker"recommended not to lower it"
28Huge cross-AZ data transfer billConsumers fetching from the leader across zonesclient.rack + broker RackAwareReplicaSelector
29Rebalance times out during onPartitionsAssigned()Expensive state loading in the callbackPreparation must return within max.poll.timeout.ms
30Corrupted state after a cooperative rebalance edge caseonPartitionsLost() not implemented → onPartitionsRevoked() ran instead, while the new owner already saved its own stateImplement onPartitionsLost() carefully, avoiding conflicts with the new owner
31Garbage/exceptions when deserializingProducer serializer ≠ consumer deserializer for that topicAvro + Schema Registry so compatibility errors surface as messages, not byte diffs
32Standalone consumer silently ignores new partitionsassign() gets no notification of partition additionsPoll partitionsFor() periodically, or bounce the app on partition changes
33Standalone consumer dies; nothing reads its partitionsassign() has no failoverUse a consumer group unless you truly need fixed assignment