Available for day contractsFrom 21st September I have availability for day and half day contracts. Please contact for more information.

Contact →
mikepreston.org

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.

ConsumersMessage BrokerProducersPublishPublishDeliverDeliverFailed MessagesAcknowledgeAcknowledgeProducer 1Producer 2Queue/TopicDead Letter QueueConsumer 1Consumer 2ConsumersMessage BrokerProducersPublishPublishDeliverDeliverFailed MessagesAcknowledgeAcknowledgeProducer 1Producer 2Queue/TopicDead Letter QueueConsumer 1Consumer 2

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.

ConsumerBrokerProducerConsumerBrokerProducerProcess MessageIf ACK lost...Process Again (Duplicate)Send MessageAcknowledge ReceiptDeliver MessageAcknowledge ProcessingRedeliver MessageAcknowledgeConsumerBrokerProducerConsumerBrokerProducerProcess MessageIf ACK lost...Process Again (Duplicate)Send MessageAcknowledge ReceiptDeliver MessageAcknowledge ProcessingRedeliver MessageAcknowledge

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.

Exactly-Once StrategyNewDuplicateReceive MessageCheck Message IDProcess MessageSkip ProcessingStore Message IDAcknowledgeExactly-Once StrategyNewDuplicateReceive MessageCheck Message IDProcess MessageSkip ProcessingStore Message IDAcknowledge

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.

Deduplication StrategiesNoYesMessageContent HashMessage IDIdempotency KeyHash StoreID StoreKey StoreExists?ProcessSkipDeduplication StrategiesNoYesMessageContent HashMessage IDIdempotency KeyHash StoreID StoreKey StoreExists?ProcessSkip

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.

Partition Strategyuser_idorder_idregionMessagePartition KeyUser PartitionOrder PartitionRegion PartitionOrdering LevelsGlobal OrderingSingle PartitionPartition OrderingMultiple PartitionsNo OrderingParallel ProcessingPartition Strategyuser_idorder_idregionMessagePartition KeyUser PartitionOrder PartitionRegion PartitionOrdering LevelsGlobal OrderingSingle PartitionPartition OrderingMultiple PartitionsNo OrderingParallel Processing

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.

DLQ FlowSuccessFailure&lt; Max>= MaxMain QueueConsumerAcknowledgeRetry CountDead Letter QueueAnalysis/ReprocessingDLQ FlowSuccessFailure&lt; Max>= MaxMain QueueConsumerAcknowledgeRetry CountDead Letter QueueAnalysis/Reprocessing

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.

Backpressure StrategiesHigh Load DetectedStrategyRate LimitingLoad SheddingScalingBufferingSlow Down ProducerDrop Low PriorityAdd ConsumersQueue OverflowBackpressure StrategiesHigh Load DetectedStrategyRate LimitingLoad SheddingScalingBufferingSlow Down ProducerDrop Low PriorityAdd ConsumersQueue Overflow

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:

  1. Event-Driven Architecture - Designing systems around event production and consumption, including event sourcing and CQRS patterns

  2. Distributed Systems Consensus - Understanding Raft, Paxos, and other consensus algorithms that underpin reliable message delivery

  3. Stream Processing Patterns - Working with Apache Kafka Streams, Apache Flink, or AWS Kinesis for real-time data processing

  4. Microservices Communication - Comparing message queues with other inter-service communication patterns like gRPC, REST, and GraphQL

  5. Observability for Distributed Systems - Implementing distributed tracing, metrics, and logging for message-driven architectures

  6. Cloud Message Services - Deep dives into AWS SQS/SNS, Google Cloud Pub/Sub, and Azure Service Bus configurations and best practices