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
- One consumer per partition for Kafka
- Use thread pools for parallel processing
- Implement message deduplication
- Handle ordering requirements carefully
- Monitor consumer lag and throughput
Practice Problems
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 & reliabilityHow 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 decompositionAnalyze 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 degradationQuiz
1. What is the difference between pull and push consumption?
2. What is offset in message consumption?
3. What happens with 'auto.offset.reset' set to 'earliest'?
4. How does Kafka achieve consumer parallelism?
5. Why use manual acknowledgment instead of auto-ack?
Flashcards
Question
Pull vs Push consumption?
Click to reveal answer
Answer
Pull: consumer polls broker (Kafka). Push: broker pushes to consumer (RabbitMQ).
Question
What is offset in messaging?
Click to reveal answer
Answer
Unique ID tracking consumer's position in partition, enabling resume from last processed message
Question
How to achieve consumer parallelism in Kafka?
Click to reveal answer
Answer
Assign each partition to a different consumer in the consumer group
Question
Auto vs manual acknowledgment?
Click to reveal answer
Answer
Auto: simple but may lose messages. Manual: ensures processing before ack, more reliable.
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.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