Learn Labs
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 = false
  • commitAsync() in the loop + commitSync() on shutdown
  • commitSync(currentOffsets) in onPartitionsRevoked()
  • auto.offset.reset = earliest (or none to fail loudly)
  • partition.assignment.strategy = CooperativeStickyAssignor
  • max.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 instance
  • session.timeout.ms = > your restart duration, < your tolerable gap-in-processing
  • onPartitionsAssigned() = load state (must finish fast)
  • onPartitionsRevoked() = commit + flush state
  • onPartitionsLost() = implemented, conflict-aware

LOW-THROUGHPUT / COST-SENSITIVE CONSUMER

  • fetch.min.bytes raised (fewer round trips, less broker load)
  • fetch.max.wait.ms set 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.bytes increased

Monitoring

SignalWhy it matters (per this chapter)
Consumer lag per partitionThe master metric. Lag vs retention = imminent silent data loss (#18)
Rebalance rate / rebalance durationDetects #4/#5; eager rebalances show as group-wide throughput dips
Commit rate + commit failure ratecommitAsync failures are invisible unless you count them in the callback
records-lag-max, records-consumed-rate, fetch-latencyFetch tuning feedback
Time between poll() calls (your own metric)The direct predictor of max.poll.interval.ms eviction
Assigned partition count per instanceDetects assignor imbalance (#6) and idle consumers (#2)
Number of live group membersStatic membership hides deaths until session.timeout.ms
__consumer_offsets healthWhere all commits actually land

Scaling

GoalLever
more consume throughputadd consumers to the group (≤ partition count) → then add partitions (breaks keyed routing! Ch. 3) → or: one consumer → queue → N worker threads
lower rebalance costCooperativeStickyAssignor + static membership
lower broker loadraise fetch.min.bytes; avoid regex subscriptions
lower cross-AZ costclient.rack + RackAwareReplicaSelector
bound memoryfetch.max.bytes (not max.partition.fetch.bytes)
bound poll-loop timemax.poll.records

"Backup" / recovery, consumer edition

The consumer's recovery story is offset manipulation:

SituationWhat 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.


On this page