Skip to content
intermediatePhase 47 · Messaging

Consumer Groups

Scale consumption with consumer groups for load balancing.

30m
0 problems
Topic Progress0%

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

  1. Match consumers to partitions for optimal parallelism
  2. Monitor consumer lag to trigger scaling
  3. Use auto-scaling based on lag metrics
  4. Test rebalancing under load
  5. Document group purposes

Practice Problems

0/3solved
Design Consumer Groups System

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 & reliability
Consumer Groups Scaling

How 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 decomposition
Consumer Groups Failure Modes

Analyze 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 degradation

Quiz

1. What is the role of the group coordinator in Kafka?

Question 1 options

2. What is the maximum number of consumers in a group?

Question 2 options

3. What is sticky partition assignment?

Question 3 options

4. When does a consumer group rebalance?

Question 4 options

5. What is consumer lag?

Question 5 options

Flashcards

Question

What is a consumer group coordinator?

Answer

A Kafka broker that manages group membership, handles join/leave, and triggers partition rebalancing

Question

Max consumers per group?

Answer

Equal to number of partitions - more consumers than partitions leaves extras idle

Question

Sticky assignment benefit?

Answer

Minimizes partition movement during rebalance, keeping existing assignments stable

Question

What triggers rebalancing?

Answer

Consumer joins or leaves the consumer group, causing partition reassignment

Question

What is consumer lag?

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

  1. Range: Consecutive partitions
  2. Round-Robin: Even distribution
  3. Sticky: Minimize movement
  4. 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