4. Kafka Consumers: Reading Data from Kafka
4.13 What actually breaks in production — Ch. 4 consolidated
| # | Symptom | Root cause | Fix |
|---|---|---|---|
| 1 | Consumer group falls behind and never catches up | Single consumer; high-latency per-record work (DB write) | Add consumers up to the partition count; beyond that, add partitions (mind Ch. 3 §9.4) |
| 2 | Added consumers, throughput unchanged | More consumers than partitions → the extras are idle | Partition count is the ceiling |
| 3 | A second application "steals" messages from the first | Both apps used the same group.id | One group.id per application |
| 4 | Rebalance storm — group never stabilizes | Processing per batch exceeds max.poll.interval.ms → evicted → rebalance → new owner also evicted | Lower max.poll.records; raise max.poll.interval.ms; move slow work off the poll thread |
| 5 | Entire group pauses for seconds on every deploy | Eager rebalance = "stop the world" | CooperativeStickyAssignor (mind the <2.3 upgrade path) |
| 6 | Consumers assigned very uneven partition counts | RangeAssignor is the default and assigns per topic independently | RoundRobinAssignor / StickyAssignor / CooperativeStickyAssignor |
| 7 | session.timeout.ms dead air after every restart | close() not called → coordinator must wait for session timeout | Shutdown hook → wakeup() → close() |
| 8 | Local cache/state rebuilt on every restart, taking minutes | Dynamic membership → new member ID → new partitions | group.instance.id (static membership) |
| 9 | After enabling static membership, partitions sit unconsumed after a crash | Static members don't leave proactively; detection waits for session.timeout.ms | Tune session.timeout.ms: high enough to survive restarts, low enough to reassign on real failure |
| 10 | group.instance.id already exists error | Two consumers with the same group.instance.id | Unique per instance |
| 11 | Duplicate processing after every crash | Autocommit — last commit up to auto.commit.interval.ms (default 5 s) old | Manual commits after processing; duplicates can be reduced but never eliminated with autocommit |
| 12 | Silent data loss when processing throws mid-batch | Autocommit + 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 |
| 13 | Records committed but never processed (loss) | commitSync() called before finishing the batch | Commit after processing |
| 14 | Committed offset went backwards; duplicates spiked | commitAsync() retried in a callback, landing an older offset after a newer one | Don't retry naively; use the monotonic sequence-number pattern |
| 15 | Last batch reprocessed after every clean shutdown | Only commitAsync() used; the final commit failed with no "next commit" to retry it | commitAsync() in the loop + commitSync() on exit |
| 16 | Batch reprocessed after every rebalance | No offsets committed before losing partitions | commitSync(currentOffsets) in onPartitionsRevoked() |
| 17 | Off-by-one: one record duplicated (or skipped) per commit | Committed record.offset() instead of record.offset()+1 | The committed offset is the next offset to read |
| 18 | A consumer down for a long weekend silently skips everything | Committed offset aged out of retention + auto.offset.reset=latest (the default) | Set earliest (accept reprocessing) or none (fail loudly); retention ≥ max outage (Ch. 2) |
| 19 | A monthly batch consumer group "forgets" everything | Group empty longer than offsets.retention.minutes (default 7 days) → offsets deleted → behaves as brand new | Raise the broker setting, or keep a member alive, or manage offsets externally |
| 20 | Broker/network overloaded by metadata, not data | Regex subscribe() on a cluster with ~30,000+ partitions — filtering is client-side, full topic list fetched at intervals | Explicit topic lists; "bandwidth used by topic metadata can be larger than the bandwidth used to send data" |
| 21 | Regex subscription fails with authorization errors | Regex requires a full describe grant on the entire cluster | Explicit lists, or grant cluster-wide describe (usually unacceptable) |
| 22 | Migrating poll(0) to poll(Duration.ofMillis(0)) broke metadata fetching | poll(long) blocked for metadata regardless of timeout; poll(Duration) does not | Move the logic into onPartitionsAssigned() |
| 23 | ConcurrentModificationException / bizarre consumer behavior | Sharing one consumer across threads, or two same-group consumers in one thread | One consumer per thread; ExecutorService or a consumer→queue→workers pattern |
| 24 | Consumer OOM after a rebalance shrinks the group | max.partition.fetch.bytes × (now-larger) partition count | Use fetch.max.bytes — an absolute cap |
| 25 | High consumer CPU on an idle topic | Tight poll loop with a short timeout and fetch.min.bytes=1 | Raise fetch.min.bytes and/or the poll timeout |
| 26 | Latency spikes to ~500 ms on a low-traffic topic | fetch.min.bytes raised, fetch.max.wait.ms left at the 500 ms default | Lower fetch.max.wait.ms to your SLA |
| 27 | Lowering request.timeout.ms made an incident worse | Disconnect/reconnect churn against an already-overloaded broker | "recommended not to lower it" |
| 28 | Huge cross-AZ data transfer bill | Consumers fetching from the leader across zones | client.rack + broker RackAwareReplicaSelector |
| 29 | Rebalance times out during onPartitionsAssigned() | Expensive state loading in the callback | Preparation must return within max.poll.timeout.ms |
| 30 | Corrupted state after a cooperative rebalance edge case | onPartitionsLost() not implemented → onPartitionsRevoked() ran instead, while the new owner already saved its own state | Implement onPartitionsLost() carefully, avoiding conflicts with the new owner |
| 31 | Garbage/exceptions when deserializing | Producer serializer ≠ consumer deserializer for that topic | Avro + Schema Registry so compatibility errors surface as messages, not byte diffs |
| 32 | Standalone consumer silently ignores new partitions | assign() gets no notification of partition additions | Poll partitionsFor() periodically, or bounce the app on partition changes |
| 33 | Standalone consumer dies; nothing reads its partitions | assign() has no failover | Use a consumer group unless you truly need fixed assignment |