Skip to content
intermediatePhase 47 · Messaging

Consumers

Build message consumers with concurrency, offset management, and idempotency.

45m
0 problems
Topic Progress0%

Message Consumption

Message Consumption

Consumers read and process messages from queues or topics.

Consumption Models

1. Pull-based (Kafka):
   Consumer polls broker for messages
   Consumer controls pace

2. Push-based (RabbitMQ):
   Broker pushes messages to consumer
   Broker controls pace

Kafka Consumer Example

Properties props = new Properties();
props.put("bootstrap.servers", "kafka:9092");
props.put("group.id", "order-processor");
props.put("enable.auto.commit", "false");
props.put("auto.offset.reset", "earliest");

KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("orders"));

try {
    while (true) {
        ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(100));
        
        for (ConsumerRecord<String, String> record : records) {
            processOrder(record.value());
            
            // Commit after processing
            consumer.commitSync();
        }
    }
} finally {
    consumer.close();
}

RabbitMQ Consumer Example

def callback(ch, method, properties, body):
    try:
        process_message(json.loads(body))
        ch.basic_ack(delivery_tag=method.delivery_tag)  # Acknowledge
    except Exception as e:
        ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)

channel.basic_consume(
    queue='orders',
    on_message_callback=callback,
    auto_ack=False  # Manual acknowledgment
)

channel.start_consuming()

Consumer Patterns

Pattern Description Use Case
Single Consumer One consumer per queue Simple processing
Competing Consumers Multiple consumers, load balanced High throughput
Consumer Group Group of consumers, partition assignment Scalable processing
Exclusive Consumer Only one consumer allowed Critical operations

Offset Management

Offset Management

What is Offset?

Partition 0: [msg0, msg1, msg2, msg3, msg4]
                          ↑
                   Committed offset = 2
                   (msg0, msg1 processed)

Commit Strategies

1. Auto-Commit:
   - Periodic background commit
   - May lose messages on failure
   - Simple but unreliable

2. Sync Commit:
   - Blocking commit after processing
   - Durable but slower

3. Async Commit:
   - Non-blocking with callback
   - Fast but may lose last batch

Implementation

// Sync commit
consumer.commitSync();

// Async commit
consumer.commitAsync((offsets, exception) -> {
    if (exception != null) {
        log.error("Commit failed", exception);
    }
});

// Commit specific offset
Map<TopicPartition, OffsetAndMetadata> offsets = new HashMap<>();
offsets.put(new TopicPartition("orders", 0), new OffsetAndMetadata(5));
consumer.commitSync(offsets);

Offset Reset Policies

When no committed offset:

1. earliest: Start from beginning
   - Reprocess all messages
   - Use for recovery

2. latest: Start from end
   - Skip existing messages
   - Use for new consumers

3. none: Throw exception
   - Require existing offset
   - Use for strict ordering

Offset Storage

Kafka:
- Offsets stored in __consumer_offsets topic
- Managed by consumer group coordinator

RabbitMQ:
- Offsets tracked by broker
- Acknowledgment-based

Concurrency

Consumer Concurrency

Concurrency Models

1. Multi-threaded Consumer:
   Single consumer, multiple threads
   [Consumer] → [Thread 1] → [Thread 2] → [Thread 3]

2. Multi-consumer:
   Multiple consumer instances
   [Consumer 1] → [Thread 1]
   [Consumer 2] → [Thread 2]
   [Consumer 3] → [Thread 3]

3. Partition-based:
   One consumer per partition (Kafka)
   P0 → C0 (Thread 0)
   P1 → C1 (Thread 1)
   P2 → C2 (Thread 2)

Implementation

import threading
from concurrent.futures import ThreadPoolExecutor

class ConcurrentConsumer:
    def __init__(self, queue_client, num_workers=10):
        self.client = queue_client
        self.num_workers = num_workers
        self.executor = ThreadPoolExecutor(max_workers=num_workers)
    
    def start(self, handler):
        """Start concurrent consumers"""
        def worker():
            while True:
                message = self.client.receive_message()
                if message:
                    self.executor.submit(handler, message)
        
        threads = []
        for _ in range(self.num_workers):
            t = threading.Thread(target=worker)
            t.daemon = True
            t.start()
            threads.append(t)
        
        return threads

# Usage
consumer = ConcurrentConsumer(sqs_client, num_workers=10)
consumer.start(process_order)

Thread Safety

import threading
from queue import Queue

class SafeConsumer:
    def __init__(self):
        self.message_queue = Queue()
        self.lock = threading.Lock()
    
    def process_message(self, message):
        # Thread-safe processing
        with self.lock:
            # Critical section
            update_shared_state(message)

Concurrency Best Practices

  1. One consumer per partition for Kafka
  2. Use thread pools for parallel processing
  3. Implement message deduplication
  4. Handle ordering requirements carefully
  5. Monitor consumer lag and throughput

Practice Problems

0/3solved
Design Consumers System

Design a scalable Consumers 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
Consumers Scaling

How would you scale Consumers 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
Consumers Failure Modes

Analyze potential failure modes for Consumers 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 difference between pull and push consumption?

Question 1 options

2. What is offset in message consumption?

Question 2 options

3. What happens with 'auto.offset.reset' set to 'earliest'?

Question 3 options

4. How does Kafka achieve consumer parallelism?

Question 4 options

5. Why use manual acknowledgment instead of auto-ack?

Question 5 options

Flashcards

Question

Pull vs Push consumption?

Answer

Pull: consumer polls broker (Kafka). Push: broker pushes to consumer (RabbitMQ).

Question

What is offset in messaging?

Answer

Unique ID tracking consumer's position in partition, enabling resume from last processed message

Question

How to achieve consumer parallelism in Kafka?

Answer

Assign each partition to a different consumer in the consumer group

Question

Auto vs manual acknowledgment?

Answer

Auto: simple but may lose messages. Manual: ensures processing before ack, more reliable.

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.Pull consumption (Kafka) gives consumer control over pace
  • 2.Offset tracking enables resuming from last position
  • 3.Kafka parallelism: one consumer per partition
  • 4.Manual acknowledgment prevents message loss
  • 5.Monitor consumer lag for processing delays

Interview Tips

  • Compare pull vs push consumption models
  • Explain offset management strategies
  • Discuss consumer group rebalancing
  • Mention consumer lag monitoring

Cheat Sheet

Cheat Sheet: Consumers

Consumption Models

  • Pull: Consumer polls (Kafka)
  • Push: Broker pushes (RabbitMQ)

Offset Management

  • Track progress with offsets
  • Sync/Async commit
  • auto.offset.reset: earliest/latest/none

Concurrency

  • One consumer per partition (Kafka)
  • Thread pools for parallel processing
  • Handle ordering requirements

Best Practices

  • Manual acknowledgment
  • Monitor consumer lag
  • Implement deduplication
  • Handle failures gracefully