Kafka Consumer Rebalancing
Rebalancing redistributes topic partitions among consumers when group membership or subscriptions change. Understanding rebalance behavior is critical for building resilient consumers.
Rebalance Triggers
Section titled “Rebalance Triggers”Membership Changes
Section titled “Membership Changes”| Trigger | Description |
|---|---|
| Consumer joins | New consumer added to group |
| Consumer leaves | Consumer calls close() |
| Consumer crashes | Consumer stops heartbeating |
| Session timeout | No heartbeat within session.timeout.ms |
| Poll timeout | No poll() within max.poll.interval.ms |
Subscription Changes
Section titled “Subscription Changes”| Trigger | Description |
|---|---|
| Topic subscription change | Consumer changes subscribed topics |
| New partitions added | Topic partition count increases |
| Topic deletion | Subscribed topic is deleted |
Coordinator Events
Section titled “Coordinator Events”| Trigger | Description |
|---|---|
| Coordinator failover | Group coordinator broker changes |
| Coordinator restart | Coordinator broker restarts |
Rebalance Protocols
Section titled “Rebalance Protocols”Eager Protocol (Legacy)
Section titled “Eager Protocol (Legacy)”All consumers stop processing during rebalance:
Characteristics:
- Stop-the-world: All partitions revoked
- Maximum disruption during rebalance
- Simple protocol logic
Cooperative Protocol (Kafka 2.4+)
Section titled “Cooperative Protocol (Kafka 2.4+)”Only affected partitions pause:
Characteristics:
- Incremental: Only moved partitions revoked
- Minimal disruption to processing
- Two-phase rebalance (may take longer)
Enabling Cooperative Rebalancing
Section titled “Enabling Cooperative Rebalancing”partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignorMigration Required
Switching from eager to cooperative requires a rolling restart strategy. Consumers cannot mix protocols within a group.
Migration Steps
Section titled “Migration Steps”-
Configure both strategies (cooperative first):
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor,org.apache.kafka.clients.consumer.RangeAssignor -
Rolling restart all consumers
-
Remove legacy strategy:
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignor -
Rolling restart again
Rebalance Phases
Section titled “Rebalance Phases”Join Phase
Section titled “Join Phase”Sync Phase
Section titled “Sync Phase”Rebalance Listener
Section titled “Rebalance Listener”Handle partition changes with ConsumerRebalanceListener:
public class MyRebalanceListener implements ConsumerRebalanceListener { private final KafkaConsumer<String, String> consumer; private final Map<TopicPartition, OffsetAndMetadata> pendingOffsets;
@Override public void onPartitionsRevoked(Collection<TopicPartition> partitions) { // Called BEFORE partitions are revoked log.info("Partitions revoked: {}", partitions);
// Commit pending work if (!pendingOffsets.isEmpty()) { consumer.commitSync(pendingOffsets); pendingOffsets.clear(); }
// Flush any buffers flushBuffers(partitions);
// Close partition-specific resources for (TopicPartition partition : partitions) { closePartitionState(partition); } }
@Override public void onPartitionsAssigned(Collection<TopicPartition> partitions) { // Called AFTER partitions are assigned log.info("Partitions assigned: {}", partitions);
// Initialize state for new partitions for (TopicPartition partition : partitions) { initializePartitionState(partition); }
// Optionally seek to specific positions // consumer.seek(partition, savedOffset); }
@Override public void onPartitionsLost(Collection<TopicPartition> partitions) { // Cooperative only: partitions lost without clean revocation log.warn("Partitions lost: {}", partitions);
// DO NOT commit - offsets may be stale // Just clean up resources for (TopicPartition partition : partitions) { closePartitionState(partition); } }}consumer.subscribe(List.of("orders"), new MyRebalanceListener(consumer, pendingOffsets));Rebalance Timing
Section titled “Rebalance Timing”Configuration
Section titled “Configuration”| Property | Default | Description |
|---|---|---|
session.timeout.ms | 45000 | Time before consumer considered dead |
heartbeat.interval.ms | 3000 | Heartbeat frequency |
max.poll.interval.ms | 300000 | Maximum time between polls |
rebalance.timeout.ms | max.poll.interval.ms | Time allowed for rebalance |
Rebalance Timeout
Section titled “Rebalance Timeout”Consumers must complete rebalance within rebalance.timeout.ms:
rebalance.timeout.ms = max.poll.interval.ms (by default)If a consumer doesn't respond within this time, it's removed from the group.
Minimizing Rebalance Impact
Section titled “Minimizing Rebalance Impact”Static Membership
Section titled “Static Membership”Reduce rebalances on planned restarts:
group.instance.id=consumer-instance-1session.timeout.ms=300000Sticky Assignment
Section titled “Sticky Assignment”Retain partition assignments across rebalances:
partition.assignment.strategy=org.apache.kafka.clients.consumer.StickyAssignorCooperative Rebalancing
Section titled “Cooperative Rebalancing”Only revoke partitions that must move:
partition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignorOptimized Processing
Section titled “Optimized Processing”Reduce processing time to avoid poll timeouts:
// Process asynchronously to avoid max.poll.interval.ms timeoutExecutorService executor = Executors.newFixedThreadPool(10);
while (running) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
if (!records.isEmpty()) { // Pause partitions during async processing consumer.pause(consumer.assignment());
CompletableFuture<Void> processing = CompletableFuture.runAsync( () -> processRecords(records), executor );
// Keep heartbeating while processing while (!processing.isDone()) { consumer.poll(Duration.ZERO); // Heartbeat only Thread.sleep(100); }
consumer.commitSync(); consumer.resume(consumer.assignment()); }}Rebalance Monitoring
Section titled “Rebalance Monitoring”Metrics
Section titled “Metrics”| Metric | Description |
|---|---|
rebalance-latency-avg | Average rebalance duration |
rebalance-latency-max | Maximum rebalance duration |
rebalance-total | Total rebalances |
rebalance-rate-per-hour | Rebalances per hour |
failed-rebalance-total | Failed rebalances |
last-rebalance-seconds-ago | Time since last rebalance |
JMX Access
Section titled “JMX Access”// Access via JMXObjectName name = new ObjectName( "kafka.consumer:type=consumer-coordinator-metrics,client-id=my-consumer");Double rebalanceRate = (Double) mbs.getAttribute(name, "rebalance-rate-per-hour");Alerting
Section titled “Alerting”| Condition | Action |
|---|---|
| Rebalance rate > 1/hour | Investigate cause |
| Failed rebalances > 0 | Check consumer health |
| Rebalance duration > 60s | Optimize consumer |
Troubleshooting
Section titled “Troubleshooting”Frequent Rebalances
Section titled “Frequent Rebalances”Symptoms:
- High
rebalance-rate-per-hour - Consumer lag increases during rebalances
- Processing gaps in metrics
Common Causes:
| Cause | Solution |
|---|---|
| Processing too slow | Increase max.poll.interval.ms or process async |
| Network instability | Increase session.timeout.ms |
| GC pauses | Tune JVM, increase session.timeout.ms |
| Consumer crashes | Fix application bugs |
Diagnostic Steps:
# Check consumer group statekafka-consumer-groups.sh --bootstrap-server localhost:9092 \ --describe --group order-processors --state
# Monitor rebalance events in logsgrep -i "rebalance\|revoked\|assigned" consumer.logStuck Rebalance
Section titled “Stuck Rebalance”Symptoms:
- Group in
PreparingRebalancestate - No progress in offset commits
- Consumers not receiving messages
Solutions:
- Check for unresponsive consumers
- Verify network connectivity
- Force remove stuck consumer (where supported):
# Remove member from group (broker-side)kafka-consumer-groups.sh --bootstrap-server localhost:9092 \ --group order-processors --delete-members \ --member consumer-1-abc123If --delete-members is unavailable, restart the stuck consumer to trigger a clean rebalance.
Rebalance Storm
Section titled “Rebalance Storm”Symptoms:
- Cascading rebalances
- Multiple rebalances in quick succession
Prevention:
# Longer timeoutssession.timeout.ms=60000max.poll.interval.ms=600000
# Static membershipgroup.instance.id=${POD_NAME}
# Cooperative rebalancingpartition.assignment.strategy=org.apache.kafka.clients.consumer.CooperativeStickyAssignorBest Practices
Section titled “Best Practices”Configuration
Section titled “Configuration”| Practice | Recommendation |
|---|---|
| Use cooperative rebalancing | CooperativeStickyAssignor for Kafka 2.4+ |
| Use static membership | Stable deployments (Kubernetes) |
| Tune timeouts | Based on processing time and network |
Implementation
Section titled “Implementation”| Practice | Recommendation |
|---|---|
| Handle rebalance events | Implement ConsumerRebalanceListener |
| Commit before revoke | Ensure progress is saved |
| Clean up resources | Close partition-specific state |
Operations
Section titled “Operations”| Practice | Recommendation |
|---|---|
| Monitor rebalance rate | Alert on unexpected increases |
| Rolling deployments | Deploy one consumer at a time |
| Test rebalance handling | Verify behavior under rebalance |
Related Documentation
Section titled “Related Documentation”- Consumer Guide - Consumer patterns
- Consumer Groups - Group coordination
- Offset Management - Offset handling
- Configuration - Configuration reference