| name | message-queue |
| description | Design message queue architectures using Kafka, RabbitMQ, or SQS. Outputs topic/queue design, producer/consumer patterns, dead letter queues, ordering guarantees, and scaling configuration. |
| argument-hint | ["use case","message volume","ordering requirements","delivery guarantees","technology choice"] |
| allowed-tools | Read, Write, Bash |
Message Queue Architecture
Design reliable, scalable messaging systems for decoupled async communication. Choose the right technology, guarantee the right semantics, and handle failures gracefully.
Technology Decision Matrix
| Requirement | Kafka | RabbitMQ | SQS | Redis Streams |
|---|
| High throughput (>100k msg/s) | ✅ Best | ⚠️ Medium | ⚠️ Medium | ✅ Good |
| Message replay / audit | ✅ Built-in | ❌ No | ❌ No | ✅ Limited |
| Complex routing | ⚠️ Limited | ✅ Best | ⚠️ Topic filter | ❌ No |
| Exactly-once delivery | ✅ Transactions | ❌ Hard | ⚠️ With dedup | ❌ No |
| Fan-out to many consumers | ✅ Consumer groups | ✅ Exchange | ✅ SNS+SQS | ✅ Groups |
| Serverless / managed | ✅ Confluent | ✅ CloudAMQP | ✅ Native AWS | ✅ ElastiCache |
| Simple queue (single consumer) | ⚠️ Overkill | ✅ Good | ✅ Best | ✅ Good |
Process
- Define message contracts — schema, versioning, payload size.
- Choose delivery semantics — at-most-once, at-least-once, exactly-once.
- Design topics/queues — naming, partitioning, routing keys.
- Plan consumer groups — parallelism, ordering constraints.
- Configure DLQ — dead letter handling, retry limits.
- Set retention — storage vs. replay requirements.
- Implement idempotent consumers — handle redelivery safely.
- Monitor lag — consumer group lag is the critical metric.
Output Format
Kafka Design
from confluent_kafka import Producer, KafkaError
from confluent_kafka.schema_registry import SchemaRegistryClient
from confluent_kafka.schema_registry.avro import AvroSerializer
import json
import uuid
import logging
from dataclasses import dataclass, asdict
from datetime import datetime, timezone
logger = logging.getLogger(__name__)
@dataclass
class OrderCreatedEvent:
event_id: str
order_id: str
user_id: str
total_cents: int
items: list[dict]
created_at: str
schema_version: str = "1.0"
class KafkaEventProducer:
def __init__(self, bootstrap_servers: str, schema_registry_url: str = None):
self._producer = Producer({
"bootstrap.servers": bootstrap_servers,
"acks": "all",
"enable.idempotence": True,
"max.in.flight.requests.per.connection": 5,
"retries": 2147483647,
: ,
: ,
: ,
: ,
})
._dlq_topic =
() -> :
envelope = {
**event,
: event.get() (uuid.uuid4()),
: datetime.now(timezone.utc).isoformat(),
}
kafka_headers = []
headers:
k, v headers.items():
kafka_headers.append((k, (v).encode()))
():
err:
logger.error(
,
extra={: topic, : key}
)
._send_to_dlq(topic, key, envelope, (err))
:
logger.debug(
)
._producer.produce(
topic=topic,
key=key.encode() key ,
value=json.dumps(envelope).encode(),
headers=kafka_headers,
on_delivery=delivery_callback
)
():
remaining = ._producer.flush(timeout=timeout)
remaining > :
TimeoutError()
():
dlq_event = {
: original_topic,
: key,
: event,
: error,
: datetime.now(timezone.utc).isoformat(),
}
._producer.produce(
topic=._dlq_topic,
key=key.encode() key ,
value=json.dumps(dlq_event).encode()
)
confluent_kafka Consumer, KafkaError, TopicPartition
signal
threading
:
():
._consumer = Consumer({
: bootstrap_servers,
: group_id,
: ,
: ,
: max_poll_interval_ms,
: ,
: ,
: ,
: ,
})
._consumer.subscribe(topics)
._running =
():
._running:
messages = ._consumer.consume(
num_messages=batch_size,
timeout=timeout_ms /
)
messages:
failed_offsets = []
msg messages:
msg.error():
msg.error().code() == KafkaError._PARTITION_EOF:
logger.error()
:
event = json.loads(msg.value())
handler(event, msg.headers() [])
Exception e:
logger.error(
,
extra={
: msg.topic(),
: msg.partition(),
: msg.offset(),
}
)
failed_offsets.append((msg.topic(), msg.partition(), msg.offset()))
failed_offsets:
._consumer.commit(asynchronous=)
():
._running =
._consumer.close()
TOPIC_CONFIG = {
: {
: ,
: ,
: * * * * ,
: ,
: ,
: ,
},
: {
: ,
: ,
: * * * * ,
},
: {
: ,
: ,
: * * * * ,
}
}
RabbitMQ Design
import pika
import json
import time
from typing import Callable
class RabbitMQSetup:
"""Exchange + queue topology for order processing."""
def __init__(self, connection_url: str):
self.connection = pika.BlockingConnection(
pika.URLParameters(connection_url)
)
self.channel = self.connection.channel()
self.channel.basic_qos(prefetch_count=1)
def setup_topology(self):
"""Declare exchanges, queues, and bindings."""
self.channel.exchange_declare(
exchange="dlx",
exchange_type="direct",
durable=True
)
self.channel.queue_declare(
queue="dead-letter",
durable=True,
arguments={"x-queue-type": "quorum"}
)
self.channel.queue_bind(
exchange="dlx",
queue="dead-letter",
routing_key="#"
)
.channel.exchange_declare(
exchange=,
exchange_type=,
durable=
)
.channel.queue_declare(
queue=,
durable=,
arguments={
: ,
: ,
: ,
: ,
}
)
.channel.queue_bind(
exchange=,
queue=,
routing_key=
)
.channel.queue_declare(
queue=,
durable=,
arguments={
: ,
: ,
}
)
.channel.queue_bind(
exchange=,
queue=,
routing_key=
)
:
():
.channel = channel
():
.channel.basic_publish(
exchange=,
routing_key=routing_key,
body=json.dumps(event).encode(),
properties=pika.BasicProperties(
delivery_mode=pika.DeliveryMode.Persistent,
content_type=,
message_id=event.get(),
timestamp=(time.time()),
priority=priority,
)
)
:
():
.channel = channel
.queue = queue
.max_retries = max_retries
():
():
retry_count = (
(properties.headers {}).get(, )
)
:
event = json.loads(body)
handler(event)
ch.basic_ack(delivery_tag=method.delivery_tag)
Exception e:
logger.error()
retry_count < .max_retries:
headers = properties.headers {}
headers[] = retry_count +
delay_ms = * ( ** retry_count)
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=)
.channel.basic_publish(
exchange=,
routing_key=method.routing_key,
body=body,
properties=pika.BasicProperties(
headers={**headers, : delay_ms},
delivery_mode=pika.DeliveryMode.Persistent,
)
)
:
ch.basic_nack(delivery_tag=method.delivery_tag, requeue=)
.channel.basic_consume(
queue=.queue,
on_message_callback=callback
)
.channel.start_consuming()
SQS + SNS (AWS)
import boto3
import json
import uuid
from datetime import datetime, timezone
class SQSConsumer:
def __init__(self, queue_url: str, max_messages: int = 10):
self.sqs = boto3.client("sqs")
self.queue_url = queue_url
self.max_messages = max_messages
def process(self, handler: callable, delete_on_success: bool = True):
while True:
response = self.sqs.receive_message(
QueueUrl=self.queue_url,
MaxNumberOfMessages=self.max_messages,
WaitTimeSeconds=20,
MessageAttributeNames=["All"],
AttributeNames=["ApproximateReceiveCount"]
)
for message in response.get("Messages", []):
receive_count = int(
message.get("Attributes", {}).get("ApproximateReceiveCount", 1)
)
try:
body = json.loads(message[])
body:
event = json.loads(body[])
:
event = body
handler(event)
delete_on_success:
.sqs.delete_message(
QueueUrl=.queue_url,
ReceiptHandle=message[]
)
Exception e:
logger.error()
():
sqs = boto3.client()
dlq = sqs.create_queue(
QueueName=,
Attributes={
: ,
: ( * * ),
}
)
dlq_arn = sqs.get_queue_attributes(
QueueUrl=dlq[],
AttributeNames=[]
)[][]
main_queue = sqs.create_queue(
QueueName=,
Attributes={
: ,
: ,
: ,
: ( * * ),
: json.dumps({
: dlq_arn,
:
})
}
)
Idempotent Consumer Pattern
import redis
from functools import wraps
class IdempotencyStore:
def __init__(self, redis_client, ttl: int = 86400):
self.redis = redis_client
self.ttl = ttl
def is_processed(self, event_id: str) -> bool:
return bool(self.redis.exists(f"processed:{event_id}"))
def mark_processed(self, event_id: str):
self.redis.setex(f"processed:{event_id}", self.ttl, "1")
def idempotent(idempotency_store: IdempotencyStore):
"""Decorator: skip duplicate events."""
def decorator(func):
@wraps(func)
def wrapper(event: dict, *args, **kwargs):
event_id = event.get("event_id")
if not event_id:
logger.warning("Event missing event_id — cannot ensure idempotency")
func(event, *args, **kwargs)
idempotency_store.is_processed(event_id):
logger.info()
result = func(event, *args, **kwargs)
idempotency_store.mark_processed(event_id)
result
wrapper
decorator
store = IdempotencyStore(redis_client)
():
order = create_order_in_db(event)
send_confirmation_email(order)
update_inventory(order)
Monitoring
groups:
- name: kafka
rules:
- alert: KafkaConsumerLagHigh
expr: kafka_consumer_group_lag > 10000
for: 5m
labels:
severity: warning
annotations:
summary: "Consumer group {{ $labels.group }} lag > 10k on {{ $labels.topic }}"
- alert: KafkaConsumerLagCritical
expr: kafka_consumer_group_lag > 100000
for: 2m
labels:
severity: critical
- alert: KafkaProducerErrors
expr: rate(kafka_producer_failed_requests_total[5m]) > 0
labels:
severity: warning
resource "aws_cloudwatch_metric_alarm" "dlq_messages" {
alarm_name = "orders-dlq-not-empty"
[]
{
}
}
Rules
- DLQ on every queue — unprocessable messages must go somewhere, never silently drop.
- Alert on DLQ depth — any message in DLQ needs human attention.
- Idempotent consumers always — at-least-once delivery means duplicates are inevitable.
- Include
event_id in every message — enables idempotency and deduplication.
- Version your event schemas — add
schema_version to every event.
- Consumer lag is the SLO — not throughput. High lag = falling behind.
- Partition by natural key — order events by
order_id, not randomly, for ordering guarantees.
- Don't share consumer groups across environments — staging and prod must never compete.
- Manual offset commit — auto-commit loses messages on crash before processing completes.
- Monitor the DLQ — don't just file and forget DLQ messages, they represent processing failures.