4. Kafka Consumers: Reading Data from Kafka
4.14 Deploy / monitor / scale / recover
The consumer's recovery story is offset manipulation:
Config recipes
AT-LEAST-ONCE, MINIMAL DUPLICATES — the common correct default
enable.auto.commit = falsecommitAsync()in the loop +commitSync()on shutdowncommitSync(currentOffsets)inonPartitionsRevoked()auto.offset.reset = earliest(ornoneto fail loudly)partition.assignment.strategy = CooperativeStickyAssignormax.poll.records= tuned so processing time ≪max.poll.interval.ms- shutdown hook →
wakeup()→close()
STATEFUL CONSUMER WITH AN EXPENSIVE LOCAL CACHE
group.instance.id= a stable unique value per instancesession.timeout.ms= > your restart duration, < your tolerable gap-in-processingonPartitionsAssigned()= load state (must finish fast)onPartitionsRevoked()= commit + flush stateonPartitionsLost()= implemented, conflict-aware
LOW-THROUGHPUT / COST-SENSITIVE CONSUMER
fetch.min.bytesraised (fewer round trips, less broker load)fetch.max.wait.msset to your latency SLA- poll timeout long (less CPU spinning on an empty topic)
MULTI-AZ / MULTI-DC CONSUMER
client.rack = <zone>- broker:
replica.selector.class = RackAwareReplicaSelector receive.buffer.bytes/send.buffer.bytesincreased
Monitoring
| Signal | Why it matters (per this chapter) |
|---|---|
| Consumer lag per partition | The master metric. Lag vs retention = imminent silent data loss (#18) |
| Rebalance rate / rebalance duration | Detects #4/#5; eager rebalances show as group-wide throughput dips |
| Commit rate + commit failure rate | commitAsync failures are invisible unless you count them in the callback |
records-lag-max, records-consumed-rate, fetch-latency | Fetch tuning feedback |
Time between poll() calls (your own metric) | The direct predictor of max.poll.interval.ms eviction |
| Assigned partition count per instance | Detects assignor imbalance (#6) and idle consumers (#2) |
| Number of live group members | Static membership hides deaths until session.timeout.ms |
__consumer_offsets health | Where all commits actually land |
Scaling
| Goal | Lever |
|---|---|
| more consume throughput | add consumers to the group (≤ partition count) → then add partitions (breaks keyed routing! Ch. 3) → or: one consumer → queue → N worker threads |
| lower rebalance cost | CooperativeStickyAssignor + static membership |
| lower broker load | raise fetch.min.bytes; avoid regex subscriptions |
| lower cross-AZ cost | client.rack + RackAwareReplicaSelector |
| bound memory | fetch.max.bytes (not max.partition.fetch.bytes) |
| bound poll-loop time | max.poll.records |
"Backup" / recovery, consumer edition
The consumer's recovery story is offset manipulation:
| Situation | What you do |
|---|---|
| Lost your downstream output file/table? | seekToBeginning() or offsetsForTimes() + seek() → replay from Kafka (bounded by broker RETENTION — Ch. 2) |
| Fell hopelessly behind on a time-sensitive stream? | seekToEnd() or seek() forward → deliberately skip |
| Corrupted consumer state? | reset the group's offsets, rebuild from the log |
Group offsets vanished (offsets.retention.minutes)? | auto.offset.reset decides. Choose it DELIBERATELY. |
The important framing: in Kafka, a consumer's "backup" is the log plus a correct offset. That's why offset-commit discipline is not a code-style concern — it is your recovery point objective.