Group Coordination
Consumer Group Coordination
Coordinator Role
Kafka Broker as Group Coordinator:
1. Manages group membership
2. Handles partition assignment
3. Processes join/leave requests
4. Triggers rebalancing
Group States
Consumer Group States:
1. Empty: No members, waiting
2. Preparing Rebalance: Collecting joins
3. Completing Rebalance: Assigning partitions
4. Stable: Normal operation
5. Dead: Group failed
State Transitions:
Empty → Preparing → Completing → Stable
↓
(member leaves)
↓
Preparing → Completing → Stable
Heartbeat Mechanism
// Heartbeat configuration
props.put("heartbeat.interval.ms", "3000");
props.put("session.timeout.ms", "10000");
props.put("max.poll.interval.ms", "300000");
// Heartbeat flow:
// 1. Consumer sends heartbeat every 3s
// 2. Broker expects heartbeat within session timeout
// 3. If no heartbeat in 10s, consumer considered dead
// 4. Group rebalances to reassign partitions
Coordinator Assignment
Partition 0: Coordinator on Broker 1
Partition 1: Coordinator on Broker 2
Partition 2: Coordinator on Broker 1
Each partition has its own coordinator broker
Group Metadata
# Consumer group metadata
group_metadata = {
'group.id': 'order-processor',
'state': 'Stable',
'members': [
{'consumer_id': 'consumer-1', 'partitions': [0, 1]},
{'consumer_id': 'consumer-2', 'partitions': [2, 3]}
],
'coordinator': 'broker-1:9092'
}
Load Balancing
Consumer Group Load Balancing
Load Balancing Strategies
1. Range Assignment:
Partitions divided into ranges
C0: P0-P2
C1: P3-P5
Simple but may be uneven
2. Round-Robin:
Distribute evenly
C0: P0, P3
C1: P1, P4
C2: P2, P5
3. Sticky:
Minimize movement during rebalance
Keep existing assignments
Only move what's necessary
4. Cooperative Sticky:
Non-stop processing
Only affected partitions revoked
Assignment Implementation
// Sticky partition assignment
props.put("partition.assignment.strategy",
"org.apache.kafka.clients.consumer.StickyAssignor");
// Or cooperative (Kafka 2.4+)
props.put("partition.assignment.strategy",
"org.apache.kafka.clients.consumer.CooperativeStickyAssignor");
Load Distribution Analysis
def analyze_load_distribution(consumer_group):
"""Analyze partition distribution across consumers"""
assignments = get_partition_assignments(consumer_group)
distribution = {}
for consumer, partitions in assignments.items():
distribution[consumer] = len(partitions)
# Calculate imbalance
values = list(distribution.values())
avg = sum(values) / len(values)
max_deviation = max(abs(v - avg) for v in values)
if max_deviation > 1:
alert(f"Load imbalance detected: {distribution}")
return distribution
When Load is Uneven
Problem:
C0: P0 (100K msgs), P1 (10K msgs) = 110K
C1: P2 (10K msgs), P3 (10K msgs) = 20K
Solutions:
1. Increase partitions (more granularity)
2. Better key distribution
3. Use sticky assignor
4. Custom partitioning strategy
Scaling Consumers
Scaling Consumers
Scaling Rules
Rule: Max parallelism = Number of partitions
3 partitions → Max 3 consumers (per group)
6 partitions → Max 6 consumers
More consumers than partitions = idle consumers
Scaling Strategies
1. Horizontal Scaling:
Add more partitions + consumers
P0, P1, P2 → C1
P3, P4, P5 → C2
2. Vertical Scaling:
More threads per consumer
[Consumer] → [Thread 1, 2, 3]
3. Multiple Consumer Groups:
Each group processes independently
Group A: P0, P1, P2
Group B: P0, P1, P2 (separate consumption)
Scaling Implementation
class AutoScaler:
def __init__(self, consumer_group, target_lag=1000):
self.group = consumer_group
self.target_lag = target_lag
def should_scale(self):
"""Check if scaling is needed"""
total_lag = self.get_total_lag()
num_consumers = self.get_consumer_count()
num_partitions = self.get_partition_count()
# Scale up if lag high and consumers < partitions
if total_lag > self.target_lag * num_consumers:
if num_consumers < num_partitions:
return 'scale_up'
# Scale down if lag low and consumers > 1
if total_lag < self.target_lag and num_consumers > 1:
return 'scale_down'
return 'no_change'
Consumer Group Use Cases
| Use Case | Description | Scaling |
|---|---|---|
| Single group | Load balanced processing | Add partitions + consumers |
| Multiple groups | Independent processing | Each group scales separately |
| Priority groups | Different priorities | More consumers for high priority |
Best Practices
- Match consumers to partitions for optimal parallelism
- Monitor consumer lag to trigger scaling
- Use auto-scaling based on lag metrics
- Test rebalancing under load
- Document group purposes
Practice Problems
Design a scalable Consumer Groups system. Cover high-level architecture, data model, and API design.
Solution
// Complete system design:
// - Functional + Non-functional requirements
// - Capacity estimation
// - Data model (SQL/NoSQL choice)
// - API endpoints
// - Component architecture
// - Scaling strategy
// - Monitoring & reliabilityHow would you scale Consumer Groups to handle 10x the current load? Identify bottlenecks and solutions.
Solution
// Scaling approach:
// 1. Load balancing
// 2. Database sharding/replication
// 3. Cache layer (Redis)
// 4. CDN for static assets
// 5. Async processing (queues)
// 6. Microservices decompositionAnalyze potential failure modes for Consumer Groups and design mitigation strategies.
Solution
// Failure mitigation:
// 1. Redundancy (multi-AZ)
// 2. Circuit breakers
// 3. Retry with backoff
// 4. Dead letter queues
// 5. Health checks
// 6. Graceful degradationQuiz
1. What is the role of the group coordinator in Kafka?
2. What is the maximum number of consumers in a group?
3. What is sticky partition assignment?
4. When does a consumer group rebalance?
5. What is consumer lag?
Flashcards
Question
What is a consumer group coordinator?
Click to reveal answer
Answer
A Kafka broker that manages group membership, handles join/leave, and triggers partition rebalancing
Question
Max consumers per group?
Click to reveal answer
Answer
Equal to number of partitions - more consumers than partitions leaves extras idle
Question
Sticky assignment benefit?
Click to reveal answer
Answer
Minimizes partition movement during rebalance, keeping existing assignments stable
Question
What triggers rebalancing?
Click to reveal answer
Answer
Consumer joins or leaves the consumer group, causing partition reassignment
Question
What is consumer lag?
Click to reveal answer
Answer
The difference between latest message offset and consumer's committed offset - measures processing delay
Revision Notes
Key Takeaways
- 1.Consumer group coordinator manages membership and rebalancing
- 2.Max parallelism equals number of partitions
- 3.Sticky assignment minimizes rebalancing impact
- 4.Cooperative rebalancing allows non-stop processing
- 5.Monitor consumer lag for scaling decisions
Interview Tips
- •Explain consumer group states and transitions
- •Compare assignment strategies (Range, Round-Robin, Sticky)
- •Discuss cooperative vs eager rebalancing
- •Know how to calculate max parallelism
Cheat Sheet
Cheat Sheet: Consumer Groups
Coordinator
- Manages group membership
- Handles partition assignment
- Triggers rebalancing
States
Empty → Preparing → Completing → Stable
Assignment Strategies
- Range: Consecutive partitions
- Round-Robin: Even distribution
- Sticky: Minimize movement
- Cooperative: Non-stop processing
Scaling
- Max parallelism = num partitions
- Monitor consumer lag
- Auto-scale based on lag
Best Practices
- Match consumers to partitions
- Use cooperative rebalancing