Message Queue Patterns
A guide to reliable messaging patterns for building robust distributed systems with guaranteed delivery and processing semantics.
Overview
Message queues enable asynchronous communication between distributed system components. Understanding delivery guarantees and processing patterns is essential for building reliable, fault-tolerant applications.
graph LR
subgraph Producers
P1[Producer 1]
P2[Producer 2]
end
subgraph Message Broker
Q[Queue/Topic]
DLQ[Dead Letter Queue]
end
subgraph Consumers
C1[Consumer 1]
C2[Consumer 2]
end
P1 -->|Publish| Q
P2 -->|Publish| Q
Q -->|Deliver| C1
Q -->|Deliver| C2
Q -->|Failed Messages| DLQ
C1 -.->|Acknowledge| Q
C2 -.->|Acknowledge| Q
Core Delivery Semantics
| Semantic | Description | Use Case |
|---|---|---|
| At-most-once | Message delivered 0 or 1 times | Metrics, logging |
| At-least-once | Message delivered 1 or more times | Most applications |
| Exactly-once | Message processed exactly 1 time | Financial transactions |
At-Least-Once Delivery
Key Concepts
At-least-once delivery ensures messages are never lost but may be delivered multiple times. This is the most common delivery guarantee offered by message brokers.
sequenceDiagram
participant P as Producer
participant B as Broker
participant C as Consumer
P->>B: Send Message
B->>P: Acknowledge Receipt
B->>C: Deliver Message
Note over C: Process Message
C->>B: Acknowledge Processing
Note over B,C: If ACK lost...
B->>C: Redeliver Message
Note over C: Process Again (Duplicate)
C->>B: Acknowledge
Common Patterns
Producer-Side Acknowledgement
# RabbitMQ example with publisher confirms
import pika
connection = pika.BlockingConnection()
channel = connection.channel()
# Enable publisher confirms
channel.confirm_delivery()
try:
channel.basic_publish(
exchange='',
routing_key='task_queue',
body='message',
properties=pika.BasicProperties(
delivery_mode=2, # Persistent message
),
mandatory=True
)
print("Message confirmed by broker")
except pika.exceptions.UnroutableError:
print("Message could not be routed")
Consumer-Side Acknowledgement
# Manual acknowledgement after processing
def callback(ch, method, properties, body):
try:
process_message(body)
# Only acknowledge after successful processing
ch.basic_ack(delivery_tag=method.delivery_tag)
except Exception as e:
# Reject and requeue on failure
ch.basic_nack(
delivery_tag=method.delivery_tag,
requeue=True
)
channel.basic_consume(
queue='task_queue',
on_message_callback=callback,
auto_ack=False # Manual acknowledgement
)
Kafka Consumer with Manual Commit
from kafka import KafkaConsumer
consumer = KafkaConsumer(
'my-topic',
bootstrap_servers=['localhost:9092'],
enable_auto_commit=False, # Manual commit
group_id='my-group'
)
for message in consumer:
try:
process_message(message.value)
# Commit offset after successful processing
consumer.commit()
except Exception as e:
# Handle failure - message will be redelivered
log_error(e)
Configuration Examples
RabbitMQ Queue Declaration
channel.queue_declare(
queue='durable_queue',
durable=True, # Queue survives broker restart
arguments={
'x-message-ttl': 86400000, # 24 hours
'x-max-length': 10000
}
)
AWS SQS Configuration
import boto3
sqs = boto3.client('sqs')
# Receive with visibility timeout
response = sqs.receive_message(
QueueUrl='https://sqs.eu-west-1.amazonaws.com/123456789/my-queue',
MaxNumberOfMessages=10,
VisibilityTimeout=300, # 5 minutes to process
WaitTimeSeconds=20 # Long polling
)
# Delete after successful processing
for message in response.get('Messages', []):
process_message(message['Body'])
sqs.delete_message(
QueueUrl=queue_url,
ReceiptHandle=message['ReceiptHandle']
)
Exactly-Once Processing
Key Concepts
Exactly-once processing ensures each message is processed exactly one time. This is typically achieved through a combination of at-least-once delivery and idempotent consumers.
graph TD
subgraph "Exactly-Once Strategy"
A[Receive Message] --> B{Check Message ID}
B -->|New| C[Process Message]
B -->|Duplicate| D[Skip Processing]
C --> E[Store Message ID]
E --> F[Acknowledge]
D --> F
end
Common Patterns
Idempotent Consumer with Database
import hashlib
from datetime import datetime, timedelta
class IdempotentProcessor:
def __init__(self, db_connection):
self.db = db_connection
def process(self, message_id, message_body):
# Check if already processed
if self._is_processed(message_id):
return {'status': 'duplicate', 'skipped': True}
# Process within transaction
with self.db.transaction():
# Record message as processed
self._mark_processed(message_id)
# Perform business logic
result = self._handle_message(message_body)
return {'status': 'processed', 'result': result}
def _is_processed(self, message_id):
query = "SELECT 1 FROM processed_messages WHERE id = %s"
return self.db.execute(query, (message_id,)).fetchone() is not None
def _mark_processed(self, message_id):
query = """
INSERT INTO processed_messages (id, processed_at)
VALUES (%s, %s)
"""
self.db.execute(query, (message_id, datetime.utcnow()))
def _handle_message(self, body):
# Business logic here
pass
Transactional Outbox Pattern
class OutboxProcessor:
"""
Ensures exactly-once publishing by writing messages
to database in same transaction as business operation.
"""
def process_order(self, order_data):
with self.db.transaction():
# Save business entity
order = Order.create(order_data)
# Write event to outbox table (same transaction)
self.db.execute("""
INSERT INTO outbox (
aggregate_id, event_type, payload, created_at
) VALUES (%s, %s, %s, %s)
""", (
order.id,
'OrderCreated',
json.dumps(order.to_dict()),
datetime.utcnow()
))
return order
def publish_outbox_messages(self):
"""Separate process publishes from outbox."""
messages = self.db.execute("""
SELECT * FROM outbox
WHERE published_at IS NULL
ORDER BY created_at
FOR UPDATE SKIP LOCKED
""").fetchall()
for msg in messages:
self.message_broker.publish(msg.event_type, msg.payload)
self.db.execute(
"UPDATE outbox SET published_at = %s WHERE id = %s",
(datetime.utcnow(), msg.id)
)
Kafka Exactly-Once Semantics
from kafka import KafkaProducer, KafkaConsumer
# Producer with transactions
producer = KafkaProducer(
bootstrap_servers=['localhost:9092'],
transactional_id='my-transactional-id',
acks='all',
enable_idempotence=True
)
producer.init_transactions()
try:
producer.begin_transaction()
producer.send('output-topic', value=b'processed-data')
producer.commit_transaction()
except Exception as e:
producer.abort_transaction()
raise
# Consumer with read-committed isolation
consumer = KafkaConsumer(
'input-topic',
bootstrap_servers=['localhost:9092'],
isolation_level='read_committed', # Only read committed messages
enable_auto_commit=False
)
Message Deduplication
Key Concepts
Message deduplication prevents duplicate messages from being processed multiple times. This is critical when at-least-once delivery may result in duplicates.
graph LR
subgraph "Deduplication Strategies"
A[Message] --> B{Content Hash}
A --> C{Message ID}
A --> D{Idempotency Key}
B --> E[Hash Store]
C --> F[ID Store]
D --> G[Key Store]
E --> H{Exists?}
F --> H
G --> H
H -->|No| I[Process]
H -->|Yes| J[Skip]
end
Common Patterns
Content-Based Deduplication
import hashlib
import redis
class ContentDeduplicator:
def __init__(self, redis_client, ttl_seconds=3600):
self.redis = redis_client
self.ttl = ttl_seconds
def is_duplicate(self, message_body):
# Generate hash from content
content_hash = hashlib.sha256(
message_body.encode()
).hexdigest()
key = f"dedup:{content_hash}"
# Check and set atomically
result = self.redis.set(
key,
'1',
nx=True, # Only set if not exists
ex=self.ttl # Expire after TTL
)
return result is None # True if key already existed
def process_if_new(self, message_body, handler):
if self.is_duplicate(message_body):
return {'status': 'duplicate', 'processed': False}
result = handler(message_body)
return {'status': 'processed', 'result': result}
ID-Based Deduplication with Sliding Window
class SlidingWindowDeduplicator:
"""
Maintains a sliding window of processed message IDs.
Memory-efficient for high-throughput systems.
"""
def __init__(self, redis_client, window_size=10000):
self.redis = redis_client
self.window_size = window_size
self.key = "dedup:window"
def check_and_add(self, message_id):
# Use sorted set with timestamp as score
timestamp = time.time()
pipe = self.redis.pipeline()
# Check if ID exists
pipe.zscore(self.key, message_id)
# Add new ID
pipe.zadd(self.key, {message_id: timestamp})
# Trim to window size
pipe.zremrangebyrank(self.key, 0, -self.window_size - 1)
results = pipe.execute()
# If score was None, ID is new
return results[0] is not None # True if duplicate
Bloom Filter Deduplication
from pybloom_live import ScalableBloomFilter
class BloomFilterDeduplicator:
"""
Probabilistic deduplication using Bloom filters.
Memory-efficient but may have false positives.
"""
def __init__(self, initial_capacity=100000, error_rate=0.001):
self.bloom = ScalableBloomFilter(
initial_capacity=initial_capacity,
error_rate=error_rate
)
def is_duplicate(self, message_id):
if message_id in self.bloom:
return True # Possibly duplicate
self.bloom.add(message_id)
return False
AWS SQS with Deduplication
import boto3
import hashlib
sqs = boto3.client('sqs')
# FIFO queue with content-based deduplication
response = sqs.send_message(
QueueUrl='https://sqs.eu-west-1.amazonaws.com/123456789/my-queue.fifo',
MessageBody='Order details...',
MessageGroupId='orders', # Required for FIFO
# Option 1: Explicit deduplication ID
MessageDeduplicationId='order-12345',
# Option 2: Content-based (enable on queue)
# ContentBasedDeduplication=True
)
Ordering Guarantees
Key Concepts
Message ordering ensures messages are processed in the sequence they were sent. Different systems provide varying levels of ordering guarantees.
graph TD
subgraph "Ordering Levels"
A[Global Ordering] --> B[Single Partition]
C[Partition Ordering] --> D[Multiple Partitions]
E[No Ordering] --> F[Parallel Processing]
end
subgraph "Partition Strategy"
G[Message] --> H{Partition Key}
H -->|user_id| I[User Partition]
H -->|order_id| J[Order Partition]
H -->|region| K[Region Partition]
end
Common Patterns
Kafka Partition-Based Ordering
from kafka import KafkaProducer, KafkaConsumer
# Producer with partition key
producer = KafkaProducer(
bootstrap_servers=['localhost:9092'],
key_serializer=str.encode,
value_serializer=lambda v: json.dumps(v).encode()
)
# Messages with same key go to same partition
def send_order_event(order_id, event_data):
producer.send(
'order-events',
key=str(order_id), # Partition key
value=event_data
)
producer.flush()
# Consumer processes partitions in order
consumer = KafkaConsumer(
'order-events',
bootstrap_servers=['localhost:9092'],
group_id='order-processor',
enable_auto_commit=False,
max_poll_records=100
)
for message in consumer:
# Messages from same partition arrive in order
process_order_event(message.key, message.value)
consumer.commit()
RabbitMQ Single Active Consumer
# Declare queue with single active consumer
channel.queue_declare(
queue='ordered_queue',
durable=True,
arguments={
'x-single-active-consumer': True # Only one consumer active
}
)
# Multiple consumers can subscribe, but only one processes
channel.basic_consume(
queue='ordered_queue',
on_message_callback=callback,
auto_ack=False
)
Sequence Number Validation
class SequenceValidator:
"""
Validates message ordering using sequence numbers.
Detects gaps and out-of-order delivery.
"""
def __init__(self, redis_client):
self.redis = redis_client
def validate_and_process(self, partition_key, sequence_num, message, handler):
last_seq_key = f"seq:{partition_key}"
with self.redis.pipeline() as pipe:
while True:
try:
pipe.watch(last_seq_key)
last_seq = pipe.get(last_seq_key)
last_seq = int(last_seq) if last_seq else 0
if sequence_num <= last_seq:
# Duplicate or old message
return {'status': 'skipped', 'reason': 'old_sequence'}
if sequence_num > last_seq + 1:
# Gap detected - may need to wait or fetch missing
return {'status': 'gap', 'expected': last_seq + 1}
# Process message
result = handler(message)
# Update sequence atomically
pipe.multi()
pipe.set(last_seq_key, sequence_num)
pipe.execute()
return {'status': 'processed', 'result': result}
except redis.WatchError:
continue # Retry on conflict
AWS SQS FIFO with Message Groups
import boto3
sqs = boto3.client('sqs')
# Send messages with ordering within groups
def send_fifo_message(queue_url, body, group_id, dedup_id):
return sqs.send_message(
QueueUrl=queue_url,
MessageBody=body,
MessageGroupId=group_id, # Orders within group maintained
MessageDeduplicationId=dedup_id
)
# Example: Order events for same customer stay ordered
send_fifo_message(
queue_url='https://sqs.eu-west-1.amazonaws.com/.../orders.fifo',
body='{"event": "order_created", "order_id": "123"}',
group_id='customer-456', # All customer-456 orders in sequence
dedup_id='order-123-created'
)
Dead Letter Queues
Key Concepts
Dead letter queues (DLQs) capture messages that fail processing after multiple attempts. They prevent poison messages from blocking queues and enable debugging and recovery.
graph LR
subgraph "DLQ Flow"
A[Main Queue] --> B[Consumer]
B -->|Success| C[Acknowledge]
B -->|Failure| D{Retry Count}
D -->|< Max| A
D -->|>= Max| E[Dead Letter Queue]
E --> F[Analysis/Reprocessing]
end
Common Patterns
RabbitMQ Dead Letter Exchange
import pika
connection = pika.BlockingConnection()
channel = connection.channel()
# Declare DLQ
channel.queue_declare(queue='dlq', durable=True)
# Declare main queue with DLX settings
channel.queue_declare(
queue='main_queue',
durable=True,
arguments={
'x-dead-letter-exchange': '',
'x-dead-letter-routing-key': 'dlq',
'x-message-ttl': 300000, # 5 minutes TTL
'x-max-length': 10000
}
)
# Consumer with retry logic
def callback(ch, method, properties, body):
retry_count = (properties.headers or {}).get('x-retry-count', 0)
try:
process_message(body)
ch.basic_ack(delivery_tag=method.delivery_tag)
except Exception as e:
if retry_count < 3:
# Republish with incremented retry count
ch.basic_publish(
exchange='',
routing_key='main_queue',
body=body,
properties=pika.BasicProperties(
headers={'x-retry-count': retry_count + 1}
)
)
ch.basic_ack(delivery_tag=method.delivery_tag)
else:
# Let it go to DLQ
ch.basic_nack(
delivery_tag=method.delivery_tag,
requeue=False
)
AWS SQS Dead Letter Queue
import boto3
import json
sqs = boto3.client('sqs')
# Create DLQ
dlq_response = sqs.create_queue(QueueName='my-dlq')
dlq_arn = sqs.get_queue_attributes(
QueueUrl=dlq_response['QueueUrl'],
AttributeNames=['QueueArn']
)['Attributes']['QueueArn']
# Create main queue with DLQ policy
redrive_policy = {
'deadLetterTargetArn': dlq_arn,
'maxReceiveCount': '3' # Move to DLQ after 3 failures
}
main_queue = sqs.create_queue(
QueueName='my-main-queue',
Attributes={
'RedrivePolicy': json.dumps(redrive_policy),
'VisibilityTimeout': '300'
}
)
# DLQ monitoring and reprocessing
def process_dlq():
while True:
response = sqs.receive_message(
QueueUrl=dlq_response['QueueUrl'],
MaxNumberOfMessages=10,
MessageAttributeNames=['All']
)
for message in response.get('Messages', []):
# Analyse failure
analyse_failure(message)
# Optionally reprocess or move to main queue
if should_retry(message):
sqs.send_message(
QueueUrl=main_queue['QueueUrl'],
MessageBody=message['Body']
)
# Delete from DLQ
sqs.delete_message(
QueueUrl=dlq_response['QueueUrl'],
ReceiptHandle=message['ReceiptHandle']
)
Kafka Error Topic Pattern
from kafka import KafkaConsumer, KafkaProducer
producer = KafkaProducer(bootstrap_servers=['localhost:9092'])
consumer = KafkaConsumer(
'main-topic',
bootstrap_servers=['localhost:9092'],
group_id='processor',
enable_auto_commit=False
)
def process_with_dlq():
for message in consumer:
retry_count = int(
dict(message.headers or []).get(b'retry-count', b'0')
)
try:
process_message(message.value)
consumer.commit()
except RecoverableError as e:
if retry_count < 3:
# Send to retry topic with delay
producer.send(
'retry-topic',
value=message.value,
headers=[
('retry-count', str(retry_count + 1).encode()),
('original-topic', b'main-topic'),
('error', str(e).encode())
]
)
else:
# Send to DLQ
producer.send(
'dlq-topic',
value=message.value,
headers=[
('error', str(e).encode()),
('failed-at', datetime.utcnow().isoformat().encode())
]
)
consumer.commit()
except UnrecoverableError as e:
# Immediately to DLQ
producer.send('dlq-topic', value=message.value)
consumer.commit()
DLQ Analysis and Metrics
import json
from datetime import datetime
from collections import Counter
class DLQAnalyser:
def __init__(self, dlq_client):
self.dlq = dlq_client
self.error_counter = Counter()
def analyse_messages(self, limit=100):
messages = self.dlq.receive_messages(limit)
analysis = {
'total': len(messages),
'errors': Counter(),
'timestamps': [],
'samples': []
}
for msg in messages:
error_type = self._extract_error_type(msg)
analysis['errors'][error_type] += 1
analysis['timestamps'].append(msg.get('timestamp'))
if len(analysis['samples']) < 5:
analysis['samples'].append({
'id': msg.get('id'),
'error': error_type,
'body_preview': msg.get('body', '')[:200]
})
return analysis
def _extract_error_type(self, message):
# Parse error from headers or body
headers = message.get('headers', {})
return headers.get('error-type', 'unknown')
Backpressure Handling
Key Concepts
Backpressure occurs when producers generate messages faster than consumers can process them. Proper handling prevents system overload and maintains stability.
graph TD
subgraph "Backpressure Strategies"
A[High Load Detected] --> B{Strategy}
B --> C[Rate Limiting]
B --> D[Load Shedding]
B --> E[Scaling]
B --> F[Buffering]
C --> G[Slow Down Producer]
D --> H[Drop Low Priority]
E --> I[Add Consumers]
F --> J[Queue Overflow]
end
Common Patterns
Producer Rate Limiting
import time
from threading import Semaphore
from functools import wraps
class RateLimitedProducer:
def __init__(self, producer, rate_per_second=100):
self.producer = producer
self.rate = rate_per_second
self.semaphore = Semaphore(rate_per_second)
self.last_reset = time.time()
def send(self, topic, message):
self._check_rate_limit()
return self.producer.send(topic, message)
def _check_rate_limit(self):
current_time = time.time()
# Reset semaphore every second
if current_time - self.last_reset >= 1:
# Release all permits
while self.semaphore._value < self.rate:
self.semaphore.release()
self.last_reset = current_time
# Acquire permit (blocks if rate exceeded)
self.semaphore.acquire()
# Token bucket rate limiter
class TokenBucket:
def __init__(self, capacity, refill_rate):
self.capacity = capacity
self.tokens = capacity
self.refill_rate = refill_rate
self.last_refill = time.time()
def acquire(self, tokens=1):
self._refill()
if self.tokens >= tokens:
self.tokens -= tokens
return True
return False
def _refill(self):
now = time.time()
elapsed = now - self.last_refill
refill_amount = elapsed * self.refill_rate
self.tokens = min(self.capacity, self.tokens + refill_amount)
self.last_refill = now
Consumer Concurrency Control
import asyncio
from concurrent.futures import ThreadPoolExecutor
class BackpressureAwareConsumer:
def __init__(self, consumer, max_concurrent=10):
self.consumer = consumer
self.semaphore = asyncio.Semaphore(max_concurrent)
self.executor = ThreadPoolExecutor(max_workers=max_concurrent)
self.queue_depth = 0
async def process_messages(self):
while True:
# Fetch based on current capacity
fetch_count = min(10, self.semaphore._value)
if fetch_count == 0:
await asyncio.sleep(0.1)
continue
messages = self.consumer.poll(
max_records=fetch_count,
timeout_ms=1000
)
for message in messages:
await self.semaphore.acquire()
asyncio.create_task(self._process_with_release(message))
async def _process_with_release(self, message):
try:
await self._process(message)
self.consumer.commit(message)
finally:
self.semaphore.release()
async def _process(self, message):
loop = asyncio.get_event_loop()
await loop.run_in_executor(
self.executor,
process_message,
message.value
)
Queue-Based Backpressure with Monitoring
import boto3
from datetime import datetime
class AdaptiveConsumer:
def __init__(self, queue_url, sqs_client=None):
self.sqs = sqs_client or boto3.client('sqs')
self.queue_url = queue_url
self.batch_size = 10
self.min_batch = 1
self.max_batch = 10
def run(self):
while True:
# Get queue metrics
depth = self._get_queue_depth()
# Adjust batch size based on queue depth
self._adjust_batch_size(depth)
# Process messages
messages = self._receive_messages()
for msg in messages:
self._process_message(msg)
# Slow down if queue is nearly empty
if depth < 10:
time.sleep(1)
def _get_queue_depth(self):
attrs = self.sqs.get_queue_attributes(
QueueUrl=self.queue_url,
AttributeNames=['ApproximateNumberOfMessages']
)
return int(attrs['Attributes']['ApproximateNumberOfMessages'])
def _adjust_batch_size(self, depth):
if depth > 1000:
self.batch_size = self.max_batch
elif depth > 100:
self.batch_size = max(self.min_batch, self.batch_size - 1)
else:
self.batch_size = self.min_batch
def _receive_messages(self):
response = self.sqs.receive_message(
QueueUrl=self.queue_url,
MaxNumberOfMessages=self.batch_size,
WaitTimeSeconds=20
)
return response.get('Messages', [])
Circuit Breaker Pattern
from enum import Enum
from datetime import datetime, timedelta
class CircuitState(Enum):
CLOSED = 'closed'
OPEN = 'open'
HALF_OPEN = 'half_open'
class CircuitBreaker:
def __init__(
self,
failure_threshold=5,
recovery_timeout=30,
half_open_requests=3
):
self.failure_threshold = failure_threshold
self.recovery_timeout = recovery_timeout
self.half_open_requests = half_open_requests
self.state = CircuitState.CLOSED
self.failure_count = 0
self.last_failure_time = None
self.half_open_successes = 0
def call(self, func, *args, **kwargs):
if self.state == CircuitState.OPEN:
if self._should_attempt_reset():
self.state = CircuitState.HALF_OPEN
self.half_open_successes = 0
else:
raise CircuitBreakerOpen("Circuit is open")
try:
result = func(*args, **kwargs)
self._on_success()
return result
except Exception as e:
self._on_failure()
raise
def _on_success(self):
if self.state == CircuitState.HALF_OPEN:
self.half_open_successes += 1
if self.half_open_successes >= self.half_open_requests:
self.state = CircuitState.CLOSED
self.failure_count = 0
else:
self.failure_count = 0
def _on_failure(self):
self.failure_count += 1
self.last_failure_time = datetime.utcnow()
if self.failure_count >= self.failure_threshold:
self.state = CircuitState.OPEN
def _should_attempt_reset(self):
if self.last_failure_time is None:
return True
elapsed = datetime.utcnow() - self.last_failure_time
return elapsed.total_seconds() >= self.recovery_timeout
Quick Reference
| Pattern | Use Case | Key Consideration |
|---|---|---|
| At-least-once | Most applications | Consumer must be idempotent |
| Exactly-once | Financial, critical ops | Higher latency, complexity |
| Content deduplication | Stateless producers | Hash collisions possible |
| ID deduplication | Known message IDs | Storage for ID tracking |
| FIFO ordering | Dependent operations | Limits parallelism |
| Partition ordering | Scalable ordering | Choose partition key wisely |
| Dead letter queue | Poison message handling | Monitor and alert on DLQ |
| Rate limiting | Protect downstream | May increase latency |
| Circuit breaker | Cascading failures | Tune thresholds carefully |
Broker Comparison
| Feature | RabbitMQ | Kafka | AWS SQS |
|---|---|---|---|
| At-least-once | Yes | Yes | Yes |
| Exactly-once | Manual | Yes (EOS) | FIFO only |
| Ordering | Per queue | Per partition | FIFO queues |
| DLQ support | Via DLX | Manual | Built-in |
| Max message size | 128MB | 1MB default | 256KB |
| Retention | Acknowledged | Time/size based | 14 days max |
Configuration Checklist
# Essential settings to configure
- [ ] Message persistence enabled
- [ ] Consumer acknowledgement mode (manual recommended)
- [ ] Visibility timeout / prefetch count
- [ ] Dead letter queue configured
- [ ] Retry policy defined
- [ ] Message TTL set
- [ ] Queue size limits
- [ ] Monitoring and alerting
Common Issues and Solutions
Issue: Duplicate Message Processing
Symptoms: Same message processed multiple times, duplicate records in database
Solutions:
# Solution 1: Idempotency key in database
CREATE UNIQUE INDEX idx_idempotency ON transactions(idempotency_key);
# Solution 2: Check-and-set with Redis
def process_if_new(message_id, handler):
if redis.set(f"processed:{message_id}", "1", nx=True, ex=3600):
return handler()
return None # Already processed
# Solution 3: Conditional database write
INSERT INTO orders (id, data)
VALUES (%s, %s)
ON CONFLICT (id) DO NOTHING;
Issue: Message Loss
Symptoms: Messages disappear without processing
Solutions:
# Enable persistence
channel.queue_declare(queue='important', durable=True)
# Use publisher confirms
channel.confirm_delivery()
# Manual acknowledgement
channel.basic_consume(queue='important', auto_ack=False)
# After successful processing only
ch.basic_ack(delivery_tag=method.delivery_tag)
Issue: Poison Messages Blocking Queue
Symptoms: Same message repeatedly fails, blocking other messages
Solutions:
# Configure retry limit and DLQ
def callback(ch, method, properties, body):
retry_count = get_retry_count(properties)
try:
process(body)
ch.basic_ack(delivery_tag=method.delivery_tag)
except Exception as e:
if retry_count >= MAX_RETRIES:
# Move to DLQ
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=False)
alert_operations(body, e)
else:
# Requeue with delay
republish_with_delay(ch, body, retry_count + 1)
ch.basic_ack(delivery_tag=method.delivery_tag)
Issue: Consumer Lag Growing
Symptoms: Messages pile up, processing falls behind
Solutions:
# 1. Increase consumer concurrency
consumer_group.scale(replicas=10)
# 2. Optimise batch processing
consumer.poll(max_records=500, timeout_ms=5000)
# 3. Add monitoring
def monitor_lag():
lag = get_consumer_lag()
if lag > THRESHOLD:
trigger_scaling_alarm()
increase_consumer_count()
# 4. Prioritise messages
# Use priority queues or multiple queues
Issue: Out-of-Order Processing
Symptoms: Events processed in wrong sequence, inconsistent state
Solutions:
# 1. Use partition keys
producer.send('topic', key=entity_id, value=event)
# 2. Include sequence numbers
message = {
'sequence': get_next_sequence(entity_id),
'data': event_data
}
# 3. Buffer and reorder
class MessageReorderer:
def process(self, message):
expected = self.get_expected_seq(message.key)
if message.seq == expected:
self.handle(message)
self.process_buffered(message.key)
else:
self.buffer(message)
Issue: Backpressure Causing Timeouts
Symptoms: Producers timeout, messages rejected
Solutions:
# 1. Implement rate limiting at producer
rate_limiter = TokenBucket(capacity=100, refill_rate=10)
def send_message(message):
while not rate_limiter.acquire():
time.sleep(0.1)
producer.send(message)
# 2. Use async sending with callbacks
def on_send_error(exc):
if isinstance(exc, BufferExhaustedError):
backoff_and_retry()
producer.send('topic', message).add_errback(on_send_error)
# 3. Monitor queue depth and scale
if queue_depth > threshold:
scale_consumers()
reduce_producer_rate()
Related Topics
The following topics would complement this Message Queue Patterns cheatsheet:
-
Event-Driven Architecture - Designing systems around event production and consumption, including event sourcing and CQRS patterns
-
Distributed Systems Consensus - Understanding Raft, Paxos, and other consensus algorithms that underpin reliable message delivery
-
Stream Processing Patterns - Working with Apache Kafka Streams, Apache Flink, or AWS Kinesis for real-time data processing
-
Microservices Communication - Comparing message queues with other inter-service communication patterns like gRPC, REST, and GraphQL
-
Observability for Distributed Systems - Implementing distributed tracing, metrics, and logging for message-driven architectures
-
Cloud Message Services - Deep dives into AWS SQS/SNS, Google Cloud Pub/Sub, and Azure Service Bus configurations and best practices