Kafka Consumers & Consumer Groups
A consumer reads records from Kafka topics, and a consumer group is a set of consumers that share the work. Kafka assigns each partition of a topic to exactly one consumer in the group, so adding consumers (up to the partition count) scales processing horizontally. Different groups read the same topic independently: the billing service and the analytics pipeline each see every event, at their own pace.
Consumers track progress by committing offsets. When and how you commit determines your delivery semantics (at-most-once or at-least-once), and rebalancing (reassigning partitions when consumers join, leave, or crash) is where most operational surprises come from.
TL;DR
- A consumer group (
group.id) divides a topic's partitions among its members: one consumer per partition at a time. - Different groups consume independently; Kafka is pub/sub across groups and a queue within a group.
- Progress is stored as committed offsets. Commit after processing for at-least-once delivery.
- Rebalances reassign partitions when membership changes; use cooperative sticky assignment and static membership to minimize disruption.
- Monitor consumer lag (latest offset minus committed offset) as the key health metric.
- Handle poison messages with retries and dead-letter topics, and make processing idempotent.
Quick Example
A Python consumer (confluent-kafka) with manual commits after processing:
Run three instances with the same group.id against a 12-partition topic, and each gets about 4 partitions.
Core Concepts
Groups and Partition Assignment
A group coordinator broker tracks membership. When the group changes, partitions are reassigned by the configured assignment strategy. Kafka 4.0's new consumer rebalance protocol (KIP-848) moves assignment to the broker and makes rebalances incremental by default.
Offsets and Commits
Each group stores its committed offset per partition in the internal __consumer_offsets topic. On restart or reassignment, a consumer resumes from the committed offset. If none exists, it uses auto.offset.reset (earliest or latest).
- Auto-commit (
enable.auto.commit=true) commits periodically in the background. It's simple, but it can commit offsets for records not yet processed (risking loss on a crash) or reprocess records after a crash. - Manual commit after processing gives at-least-once delivery: a crash before the commit means records are redelivered, never lost.
- Commit before processing gives at-most-once: a crash can skip records.
- Exactly-once requires transactions or idempotent sinks. See Kafka exactly-once.
Rebalancing
A rebalance happens when a consumer joins, leaves, crashes (misses heartbeats past session.timeout.ms), or takes too long between polls (max.poll.interval.ms). With the classic eager protocol, all consumers stop and give up their partitions during a rebalance, a "stop-the-world" pause. Mitigations:
- Cooperative sticky assignment: only the partitions that need to move are revoked, and other consumers keep working.
- Static membership (
group.instance.id): a restarting consumer rejoins with the same identity withinsession.timeout.mswithout triggering a rebalance, which is great for rolling deploys on Kubernetes. - Rebalance listeners: commit offsets and flush state in
onPartitionsRevokedto avoid reprocessing.
The Poll Loop and Processing Time
Consumers must call poll() regularly. If processing a batch takes longer than max.poll.interval.ms, the consumer is considered stuck and kicked out of the group, and its partitions get reassigned, often causing the same slow batch to be reprocessed elsewhere in a loop. Fix it by reducing max.poll.records, speeding up processing, raising the interval, or handing work to a bounded worker pool while pausing partitions.
Consumer Lag
Lag is the difference between a partition's latest offset and the group's committed offset: how far behind the consumer is. Rising lag means consumers can't keep up (scale out, up to the partition count, or speed up processing). Monitor it with Burrow, Kafka exporter plus Prometheus, or managed service dashboards, and use it to drive autoscaling with KEDA's Kafka scaler. See Kubernetes autoscaling.
Error Handling
A single bad record (malformed data, a bug-triggering payload) can block a partition forever if the consumer keeps retrying it: a poison pill. A common strategy:
- Retry transient failures a few times with backoff (network, downstream timeouts).
- On persistent failure, publish the record, with error details in headers, to a dead-letter topic (DLT) and move on.
- Optionally use retry topics with increasing delays (
orders.retry.1m,orders.retry.10m) for failures that may recover. - Alert on DLT volume, and build tooling to inspect and replay dead-lettered records.
Deserialization failures should also go to a DLT rather than crash the consumer.
Best Practices
Make Processing Idempotent
At-least-once delivery means duplicates happen during rebalances, restarts, and retries. Deduplicate by event ID, use upserts, or record processed offsets transactionally with your results. See idempotency.
Commit Deliberately
Commit after the side effects are durable. Batch commits for throughput, but not so rarely that a crash causes a huge replay. Commit synchronously on shutdown and in revocation callbacks.
Size Consumer Counts to Partitions
Running more consumers than partitions wastes resources. If you need more parallelism than partitions allow, increase partitions (with care for keyed ordering) or process records concurrently within a consumer while preserving per-key order.
Name Groups by Purpose
group.id identifies a logical subscriber (fulfillment-service), not an instance. Changing it resets progress to auto.offset.reset, which is a common accidental "reprocess everything" or "skip everything" event.
Common Mistakes
Auto-Commit With Asynchronous Processing
Commit only offsets whose processing has completed.
Slow Processing Triggering Rebalance Storms
Calling a slow API for every record in a 500-record poll batch exceeds max.poll.interval.ms, the consumer gets evicted, and its partitions bounce between instances. Lower max.poll.records, add concurrency, or raise the interval.
Infinite Retries on Poison Messages
Retrying a malformed record forever stalls the whole partition while lag grows. Cap retries and dead-letter the record.
FAQ
What happens if there are more consumers than partitions?
The extra consumers are assigned nothing and stay idle, though they act as hot standbys that take over if another consumer fails. To use more consumers, the topic needs more partitions.
How do I reprocess data from the beginning?
Reset the group's offsets with kafka-consumer-groups.sh --reset-offsets --to-earliest (or to a timestamp or specific offset) while the group is stopped, or consume with a new group.id and auto.offset.reset=earliest. Make sure processing is idempotent before replaying.
What's the difference between session.timeout.ms and max.poll.interval.ms?
session.timeout.ms detects dead consumers: the heartbeat thread stopped, usually because the process died or the network failed. max.poll.interval.ms detects stuck consumers: the process is alive but hasn't called poll() within the interval, usually because processing is too slow. Either one triggers a rebalance.
Can two groups read the same topic?
Yes, and that's the core of Kafka's pub/sub model. Each group has its own committed offsets and receives every record. Within a group, each record goes to one consumer.
Related Topics
- Kafka — The platform overview
- Kafka Topics & Partitions — What gets assigned to consumers
- Kafka Producers — Keys and ordering at the source
- Kafka Exactly-Once — Transactional consume-transform-produce
- Message Queues — Queue semantics compared
- Idempotency — Safe handling of redelivered records