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
| Model | Delivery | Example |
|---|---|---|
| Queue (work) | Each message to one of N competing consumers | SQS, RabbitMQ queue |
| Pub/Sub | Each message to every subscriber | SNS, RabbitMQ fanout |
| Log (stream) | Durable ordered log; consumers track an offset | Kafka, 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
| Guarantee | Meaning | Risk |
|---|---|---|
| At-most-once | Fire and forget; no retry | May lose messages |
| At-least-once | Retry until ack | Duplicates possible |
| Exactly-once | Processed once, no dup effect | Hard/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.
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.
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:
| Concept | Meaning |
|---|---|
| Topic | Named stream of records |
| Partition | Ordered, append-only shard of a topic; unit of parallelism |
| Offset | Position of a record in a partition; consumers commit it |
| Consumer group | Consumers share partitions; one partition → one member |
| Replication | Each partition has a leader + follower replicas (ISR) |
| Log compaction | Retains 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
| System | Model | Strength | Watch out |
|---|---|---|---|
| Kafka | Partitioned log | Very high throughput, replay, streaming | Heavier ops; no per-message TTL routing |
| RabbitMQ | Broker w/ exchanges + queues | Flexible routing, priorities, per-msg TTL | Lower throughput than Kafka at scale |
| SQS / SNS | Managed queue / pub-sub | Zero ops, auto-scale, built-in DLQ | No replay (std); FIFO throughput capped |
| NATS (JetStream) | Lightweight pub/sub + stream | Ultra-low latency, tiny footprint | Smaller ecosystem |
| Rough numbers (single cluster) | Throughput | Latency |
|---|---|---|
| Kafka | millions msg/s | ~2–10 ms |
| RabbitMQ | tens of thousands msg/s | sub-ms–ms |
| SQS | near-unbounded (std) | ~10–100 ms |
| NATS | millions msg/s | microseconds–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.