“What is a consumer group?” sounds like a warm-up question, but interviewers use it to reach the topics that hurt in production: “what triggers a rebalance?”, “why does our group rebalance every few minutes?”, “what is the difference between eager and cooperative rebalancing?” and “how do you deploy consumers on Kubernetes without a rebalance storm?” A candidate who can explain the group protocol and its timeouts usually also knows how to debug growing lag.
This guide covers how a group divides partitions, what the coordinator does, every common rebalance trigger and the tools that make rebalances cheap. It targets Kafka 3.8/4.0 with the Java client. Code is illustrative: it follows the documented APIs but was not run against a live cluster.
Before you start
You should know that a topic is split into partitions and that consumers track their progress as committed offsets per partition. If those terms are new, read Kafka partitions, offsets and ordering first. Basic Java and the idea of a heartbeat (a periodic “I am alive” message) are enough for the rest.
The short answer
A consumer group is a set of consumers with the same group.id that share the work of reading a topic: each partition is assigned to exactly one member, so every record is processed once by the group, and separate groups each get all the data. When membership changes (a member joins, leaves, misses heartbeats or stops calling poll() for longer than max.poll.interval.ms) or the topic gains partitions, the group rebalances and redistributes partitions. Eager rebalancing stops the whole group; cooperative rebalancing and the Kafka 4.0 consumer protocol move only the partitions that change, and static membership avoids rebalances for quick restarts.
How it works
Each group has a group coordinator, a broker chosen by hashing the group.id onto a partition of the internal __consumer_offsets topic. The coordinator tracks members, receives heartbeats and stores committed offsets.
With the classic protocol (still the default in 4.0), a rebalance runs in two rounds. Every member sends JoinGroup; the coordinator picks one member as the group leader and sends it the member list. The leader runs the client-side assignor and returns the plan in SyncGroup, and the coordinator hands each member its partitions. Between those rounds the group works on a shared generation number, and a member that misses the round is out.
Two timers decide when a member is considered gone, and they measure different things:
session.timeout.ms(45 seconds by default): heartbeats are sent by a background thread everyheartbeat.interval.ms(3 seconds). Missing them for the session timeout means the process is dead or frozen, for example in a long GC pause.max.poll.interval.ms(5 minutes): the maximum time between twopoll()calls. If your processing loop takes longer, the background thread itself makes the consumer leave the group, because a live process that is not polling is stuck.
The assignor decides how partitions are spread. The Java default is [RangeAssignor, CooperativeStickyAssignor]: range is used until every member supports the cooperative one, which makes upgrading easy.
| Style | Assignors | What happens on rebalance |
|---|---|---|
| Eager | RangeAssignor, RoundRobinAssignor, StickyAssignor |
every member revokes all partitions, then receives a new set |
| Cooperative | CooperativeStickyAssignor |
members keep what they still own; only moving partitions are revoked, over two quick rounds |
| Consumer protocol (KIP-848, GA in 4.0) | server-side uniform or range |
the broker computes the target and each member converges incrementally, with no group-wide barrier |
Step-by-step walkthrough
Step 1: Write a poll loop that respects the timers
The consumer is not thread-safe, so one thread owns it and calls poll() regularly. Bounding the batch size is the simplest way to stay inside max.poll.interval.ms.
Properties props = new Properties();
props.put(ConsumerConfig.BOOTSTRAP_SERVERS_CONFIG, "localhost:9092");
props.put(ConsumerConfig.GROUP_ID_CONFIG, "billing");
props.put(ConsumerConfig.ENABLE_AUTO_COMMIT_CONFIG, false);
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 100); // default 500
props.put(ConsumerConfig.KEY_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
props.put(ConsumerConfig.VALUE_DESERIALIZER_CLASS_CONFIG, StringDeserializer.class);
try (KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props)) {
consumer.subscribe(List.of("orders"), new SaveOnRevoke(consumer));
while (running) {
ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(500));
for (ConsumerRecord<String, String> r : records) handle(r); // keep this well under 5 minutes
consumer.commitAsync();
}
}If each record can take 2 seconds, 500 records per poll is 1,000 seconds, far beyond the 5-minute limit. With 100 records the worst case is about 200 seconds.
Step 2: Save progress before partitions move
A ConsumerRebalanceListener is called during the rebalance on the polling thread. Committing in onPartitionsRevoked means the next owner starts exactly where this member stopped, instead of replaying everything since the last periodic commit.
class SaveOnRevoke implements ConsumerRebalanceListener {
private final KafkaConsumer<String, String> consumer;
SaveOnRevoke(KafkaConsumer<String, String> consumer) { this.consumer = consumer; }
@Override public void onPartitionsRevoked(Collection<TopicPartition> partitions) {
consumer.commitSync(); // still a member, so the commit is accepted
}
@Override public void onPartitionsAssigned(Collection<TopicPartition> partitions) {
log.info("assigned {}", partitions);
}
@Override public void onPartitionsLost(Collection<TopicPartition> partitions) {
log.warn("lost {} without a clean handover", partitions); // do not commit: no longer the owner
}
}onPartitionsLost fires when the member was already kicked out, for example after a session timeout. Committing there would fail, and another member may already own those partitions.
Step 3: Switch to cooperative assignment
With eager assignors, adding one consumer to a 20-member group pauses all 20. Cooperative assignment pauses only the partitions that move. Because the default list already includes the cooperative assignor, one rolling restart that removes RangeAssignor completes the switch.
props.put(ConsumerConfig.PARTITION_ASSIGNMENT_STRATEGY_CONFIG,
List.of(CooperativeStickyAssignor.class));If a group explicitly uses only an eager assignor, migrate in two rolling restarts: first list both assignors, then remove the eager one. Jumping straight to cooperative-only with a mixed group fails, because members cannot agree on a protocol.
Step 4: Give pods a stable identity
In a rolling deploy, each restarted pod normally leaves and rejoins as a new member, causing two rebalances per pod. Static membership gives each instance a fixed group.instance.id, so a restart within the session timeout gets its old partitions back with no rebalance.
String pod = System.getenv("HOSTNAME"); // e.g. billing-0 in a StatefulSet
props.put(ConsumerConfig.GROUP_INSTANCE_ID_CONFIG, pod);
props.put(ConsumerConfig.SESSION_TIMEOUT_MS_CONFIG, 90_000); // cover a normal restartThe cost is slower failure detection: if that pod dies for good, its partitions wait for the full session timeout. Instance IDs must be unique; a second live consumer with the same ID fences the first.
Step 5: Consider the Kafka 4.0 consumer protocol
KIP-848 moves assignment to the broker. Each member heartbeats, the coordinator computes the target assignment, and members revoke and acquire partitions individually without stopping the group. It is opt-in on the client:
props.put(ConsumerConfig.GROUP_PROTOCOL_CONFIG, "consumer");
props.put(ConsumerConfig.GROUP_REMOTE_ASSIGNOR_CONFIG, "uniform"); // or "range"With this protocol, the client-side partition.assignment.strategy, session.timeout.ms and heartbeat.interval.ms no longer apply; the broker’s group.consumer.* settings govern timeouts.
Worked scenario
An invoicing service enriches each order by calling a tax API, then writes to PostgreSQL. After a tax provider incident, the API slows from 50 ms to 1.5 seconds per call. Lag climbs, logs show CommitFailedException, and the group rebalances every few minutes.
The arithmetic explains it: 500 records per poll at 1.5 seconds each is 750 seconds, so every member exceeds the 5-minute max.poll.interval.ms. Each one leaves the group, its partitions move, its later commit fails because it is no longer a member, and the new owner reprocesses the same batch and hits the same wall.
// Broken: unbounded batch, no timeout on the external call
ConsumerRecords<String, String> records = consumer.poll(Duration.ofSeconds(1));
for (var r : records) {
TaxQuote q = taxClient.quote(r.value()); // can block for seconds
repository.save(enrich(r.value(), q));
}
consumer.commitSync();The fix combines a smaller batch, a hard timeout on the dependency and a failure path that does not stall the poll loop:
props.put(ConsumerConfig.MAX_POLL_RECORDS_CONFIG, 50);
HttpClient http = HttpClient.newBuilder().connectTimeout(Duration.ofSeconds(1)).build();
// each request also sets .timeout(Duration.ofSeconds(2)); failures go to a retry topicFor work that is legitimately long, consumer.pause(assignment) stops fetching while you keep calling poll(), which keeps the member alive; resume() restarts fetching when the work completes.
Common mistake
- “Heartbeats keep a slow consumer in the group.” They cover
session.timeout.msonly. A consumer that does not poll withinmax.poll.interval.msleaves regardless. - “Add consumers to fix lag.” Only up to the partition count; extra members sit idle.
- “Just raise
max.poll.interval.msto an hour.” That hides a stuck consumer for an hour. Fix batch size and timeouts first. - Committing in
onPartitionsLost. The partitions already belong to someone else. - Reusing one
group.idfor unrelated services. They split the partitions between them and each sees only part of the data.
Verify the behavior
Watch the assignment change while you start and stop consumers:
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group billing --members
# CONSUMER-ID, HOST, CLIENT-ID, #PARTITIONS for each member
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group billing --state
# COORDINATOR, ASSIGNMENT-STRATEGY, STATE (Stable, PreparingRebalance, ...), #MEMBERS
kafka-consumer-groups.sh --bootstrap-server localhost:9092 --describe --group billing
# per partition: CURRENT-OFFSET, LOG-END-OFFSET, LAG, CONSUMER-IDStart a third consumer on a 2-partition topic and --members shows one member with 0 partitions. In the application, the consumer metrics rebalance-rate-per-hour, last-rebalance-seconds-ago and failed-rebalance-rate-per-hour reveal churn; alert when rebalances happen outside deploys.
Follow-up questions
How is session.timeout.ms different from max.poll.interval.ms? The first detects a dead or frozen process through missed heartbeats; the second detects a live process whose processing loop is stuck or too slow.
What happens when a topic gains partitions? Subscribed groups notice through metadata refresh and rebalance to assign the new partitions. Consumers using auto.offset.reset=latest may skip records written to them before assignment.
How do you scale beyond the partition count? Add partitions (carefully, for keyed topics) or parallelize inside each consumer by key while committing only fully processed offsets.
Why can a rebalance cause duplicates? Records processed since the last commit are processed again by the new owner, which is why commit-on-revoke and idempotent handlers matter.
Interview exercise
A group of 4 consumers reads a 12-partition topic with the default classic protocol and RangeAssignor. Records take up to 1 second each, max.poll.records is 500, and the team deploys with a rolling restart. They see lag spikes during every deploy and occasional CommitFailedException at night. Diagnose both problems and propose changes.
Answer and reasoning
The night-time errors come from the poll loop: 500 records at up to 1 second is up to 500 seconds, beyond the 300-second max.poll.interval.ms, so slow batches evict members and their commits fail. Lower max.poll.records to around 100 and put timeouts on slow calls. The deploy spikes come from eager rebalancing: each pod restart causes a leave and a join, and with RangeAssignor every member revokes all 12 partitions both times. Switching to CooperativeStickyAssignor limits each rebalance to the moving partitions, and static membership with a session timeout longer than a pod restart removes the rebalances entirely. Committing in onPartitionsRevoked reduces the duplicates that remain.
Continue learning
Practise with the Apache Kafka interview questions and the Kafka MCQs. Related notes: Kafka offset commits and delivery semantics and Kafka partitions, offsets and ordering. Primary sources: the Apache Kafka documentation, the ConsumerRebalanceListener Javadoc, KIP-848 on the new rebalance protocol and Confluent’s consumer design notes.