Learn/system design/Message Queues
Intermediate~25 min read

Message Queues

Async messaging, delivery semantics, ordering, backpressure, and the Kafka vs RabbitMQ vs SQS landscape.

KafkaDeliveryIdempotencyDLQ

Why Async Messaging?

A message queue (or broker) sits between producers and consumers so they don't call each other directly. Instead of a synchronous request that blocks until the callee responds, the producer hands a message to the broker and moves on; consumers process it later. This buys three things: decoupling (services don't know about each other), buffering (traffic spikes are absorbed rather than dropped), and resilience (a consumer can be down and catch up later).

Producer Broker Consumer(s) --------- ---------- ------------ [order svc] --msg--> [ queue / topic ] --pull/push--> [ email svc ] (durable buffer) [ ledger svc ] [ analytics ] fire & forget stores, orders, process at their low latency retries, fans out own pace

Sync vs async

Use synchronous calls when the caller needs the result now (a price, an auth check). Use messaging for work that can happen later (send email, update search index, run billing) or must fan out to many consumers. The cost is added latency and eventual consistency.

Queue vs Pub/Sub vs Log

ModelDeliveryExample
Queue (work)Each message to one of N competing consumersSQS, RabbitMQ queue
Pub/SubEach message to every subscriberSNS, RabbitMQ fanout
Log (stream)Durable ordered log; consumers track an offsetKafka, Kinesis, Pulsar

The key difference of a log: messages are not deleted on consume. They persist for a retention window, so multiple independent consumers can read the same stream at their own pace and you can replay from any offset. A classic queue deletes a message once acknowledged.

Delivery Semantics

GuaranteeMeaningRisk
At-most-onceFire and forget; no retryMay lose messages
At-least-onceRetry until ackDuplicates possible
Exactly-onceProcessed once, no dup effectHard/costly; often "effectively once"

At-least-once is the practical default. True end-to-end exactly-once across a network is impossible in general (the ack can be lost after the work is done), so systems approximate it with idempotency or transactional processing. Kafka offers exactly-once within Kafka (idempotent producer + transactions), but the moment a consumer writes to an external system, you are back to needing idempotency.

Idempotency — the real fix

Since duplicates are unavoidable with at-least-once, make processing idempotent: applying the same message twice has the same effect as applying it once. Use a unique message/idempotency key and dedupe.

python
def handle(msg):
    if seen.contains(msg.id):     # dedupe store (Redis/DB w/ TTL)
        return                    # already processed -> skip
    with db.transaction():
        apply_effect(msg)         # e.g. UPSERT, not blind INSERT
        seen.add(msg.id)          # record processed id atomically

Ordering Guarantees

Global total ordering across a distributed queue is expensive and limits throughput. Most systems offer partition-level (or per-key) ordering: messages with the same key go to the same partition and are delivered in order there, while different partitions are processed in parallel. Choose a partition key (e.g. user_id) so related events stay ordered.

Ordering vs parallelism trade-off

Strict ordering means one consumer per partition — you can only parallelize up to the partition count. Retries also threaten order: reprocessing a failed message after later ones violates it. If order matters, keep failures on the same partition or design handlers to tolerate reordering.

Backpressure & Flow Control

If producers outpace consumers, the queue grows without bound — latency balloons and the broker may run out of disk. Backpressure is the feedback that slows producers or sheds load. Tactics: bounded queues that reject/block on full, consumer-driven pull (Kafka) so consumers set the pace, autoscaling consumers on queue depth/lag, and rate limiting producers. Always alarm on consumer lag (unprocessed backlog), the single most important queue metric.

Retries, DLQs, and Poison Messages

Transient failures (a downstream timeout) deserve a retry; permanent failures (malformed payload — a poison message) will fail forever and can block the queue if retried endlessly. The pattern: retry a bounded number of times with exponential backoff + jitter, then route the message to a dead-letter queue (DLQ) for inspection.

attempt 1: fail -> wait ~1s (base * 2^0 + rand jitter) attempt 2: fail -> wait ~2s attempt 3: fail -> wait ~4s attempt 4: fail -> wait ~8s attempt 5: fail -> route to DLQ (alert, inspect, replay later) delay = min(cap, base * 2^attempt) + random_jitter

Why jitter

Without jitter, many failed messages retry at the same instant (a synchronized thundering herd) and re-overwhelm the recovering downstream. Randomizing the delay spreads retries out.

Consumption Patterns

Competing consumers (scale out a queue): [queue] --> C1 --> C2 each msg to exactly ONE consumer; --> C3 add consumers to raise throughput Fan-out (pub/sub): /--> [email topic sub] [event] --+--> [search index sub] each subscriber gets a COPY \--> [analytics sub]

Competing consumers parallelize a workload: many consumers pull from one queue, each message handled once, throughput scales with consumer count (up to partition count for logs). Fan-out delivers each event to multiple independent subscribers — the backbone of event-driven architecture, where services react to events instead of being commanded.

The Outbox Pattern

A common bug: a service writes to its DB and publishes an event as two separate steps. If it crashes between them, state and event diverge (the "dual-write" problem). The transactional outbox fixes this: write the event to an outbox table in the same DB transaction as the state change, then a relay (or CDC) reads the outbox and publishes to the broker.

sql
BEGIN;
  UPDATE orders SET status = 'PAID' WHERE id = 42;
  INSERT INTO outbox (id, topic, payload) VALUES (...);  -- same txn
COMMIT;                          -- both or neither

-- relay / CDC (Debezium) tails the outbox table -> Kafka
-- at-least-once publish; consumers dedupe by event id

Kafka Internals

Kafka is a distributed, replicated commit log. Key concepts:

ConceptMeaning
TopicNamed stream of records
PartitionOrdered, append-only shard of a topic; unit of parallelism
OffsetPosition of a record in a partition; consumers commit it
Consumer groupConsumers share partitions; one partition → one member
ReplicationEach partition has a leader + follower replicas (ISR)
Log compactionRetains only the latest value per key (changelog/state)

Because a consumer group assigns each partition to exactly one member, maximum consumer parallelism equals the partition count — choose partitions with future scale in mind. Offsets let a slow or restarted consumer resume exactly where it left off, and let you replay history.

Broker Comparison

SystemModelStrengthWatch out
KafkaPartitioned logVery high throughput, replay, streamingHeavier ops; no per-message TTL routing
RabbitMQBroker w/ exchanges + queuesFlexible routing, priorities, per-msg TTLLower throughput than Kafka at scale
SQS / SNSManaged queue / pub-subZero ops, auto-scale, built-in DLQNo replay (std); FIFO throughput capped
NATS (JetStream)Lightweight pub/sub + streamUltra-low latency, tiny footprintSmaller ecosystem
Rough numbers (single cluster)ThroughputLatency
Kafkamillions msg/s~2–10 ms
RabbitMQtens of thousands msg/ssub-ms–ms
SQSnear-unbounded (std)~10–100 ms
NATSmillions msg/smicroseconds–ms

Capacity sketch

1 KB messages at 200k msg/s = ~200 MB/s ingest. With 7-day retention that is ~120 TB before replication; at replication factor 3, ~360 TB of disk. If one partition sustains ~10k msg/s, you need at least ~20 partitions for throughput — and thus up to ~20 consumers per group. Plan partitions for peak plus growth, since increasing them later breaks key-based ordering.

Design checklist

Assume at-least-once → make consumers idempotent. Set a DLQ + bounded retries with backoff+jitter. Pick a partition key for ordering. Alarm on consumer lag. Use the outbox pattern to avoid dual-write bugs. Pick Kafka for streams/replay, RabbitMQ for rich routing, SQS/SNS for zero-ops managed queues, NATS for ultra-low latency.

Section navigation