Message Queues Explained: Kafka vs. RabbitMQ vs. SQS
"I'll put a message queue between them" is the most under-examined sentence in system design interviews. It sounds like an architecture decision, but on its own it commits you to nothing — and the choices hiding underneath it decide whether your system loses writes, reorders a user's actions, duplicates a payment, or turns a downstream outage into a four-hour backlog.
This post covers what a queue actually buys you and what it costs, why queues and logs are different products that people call the same thing, the three delivery guarantees and which one you should design for, how partitioning trades ordering against throughput, the dual-write problem that a queue quietly introduces, and how to choose between Kafka, RabbitMQ and SQS with a reason you can defend.
What a queue actually buys you
Take a video upload endpoint — the front door of the pipeline in design YouTube — that, before responding, transcodes the file, sends a confirmation email, and writes an analytics row. Every one of those is work the user does not need to wait for, and every one is a way for the upload to fail.
The arithmetic is worth doing out loud, because it is the actual argument. Four components in series, each at 99.9% availability, gives you 0.999⁴ ≈ 99.6% — about 35 hours of downtime a year instead of 8.8. Serial dependencies multiply, and every one you add makes the endpoint worse than its worst component. Publishing to one broker replaces three of those multiplications with one.
So a queue buys three specific things:
- Latency isolation. The endpoint's p99 stops being the sum of everything downstream.
- Failure isolation. A consumer being down becomes a growing backlog rather than a user-visible error, provided the work is genuinely deferrable.
- Burst absorption. Producers and consumers no longer have to be provisioned for the same instantaneous rate.
And it costs you four: an extra system to operate, eventual consistency the user can observe ("I uploaded it, where is it?"), duplicate deliveries you must now handle, and debugging that has become much harder because the causal chain is no longer a stack trace.
Queues and logs are not the same thing
This is the distinction that separates candidates who have used a broker from candidates who have read about one. "Message queue" gets applied to two products with genuinely different semantics.
A classic queue (RabbitMQ, SQS, Celery on Redis, ActiveMQ) is a work distribution system. Messages are handed to competing consumers, acked, and deleted. Its strengths are per-message acknowledgement, rich routing, per-message TTLs and priorities, and consumers that can ack out of order.
A log (Kafka, Pulsar, Kinesis, Redis Streams) is an append-only, retained, ordered sequence. Consumers track an offset; the broker keeps data for a retention window — Kafka's default is seven days — regardless of who has read it. That single property unlocks things a queue cannot do: adding a second independent consumer group later, replaying last Tuesday after you fix a bug, and rebuilding a derived store from scratch.
The practical rule I use: if the message is a command for one worker, use a queue; if it is a fact that several systems may care about now or later, use a log. "Send this email" is a command. "Order 1234 was placed" is a fact. Systems accumulate more facts than commands over time, which is why teams that start on RabbitMQ often end up running Kafka alongside it.
Delivery guarantees, and the only one you should design for
Three guarantees get named. Only one is a real engineering choice.
At-most-once means you ack before processing. A crash mid-work loses the message silently. This is defensible only for data where loss is genuinely cheap — a sampled metric, a cache-warm hint.
At-least-once means you ack after the work commits. If you crash after committing but before acking, the broker redelivers and the work happens twice. This is what every production system actually runs.
Exactly-once is not achievable as a delivery property across a network — two generals, and the acknowledgement can always be the message that gets lost. What Kafka's enable.idempotence and transactions give you is exactly-once processing within a closed read-process-write loop over Kafka topics, with the consumer offset committed in the same transaction as the output. The moment your side effect leaves that boundary — an HTTP call to Stripe, a row in Postgres — you are back to at-least-once, and you handle it in the consumer.
So the design is: at-least-once delivery plus an idempotent consumer. The mechanism is a natural or supplied idempotency key and a uniqueness constraint that the database enforces for you.
def handle(msg):
# The broker may deliver this more than once. Make the *effect* happen once.
with db.transaction():
applied = db.execute(
"INSERT INTO processed_events (event_id) VALUES (%s) "
"ON CONFLICT (event_id) DO NOTHING",
msg.id,
).rowcount
if applied == 0:
return # a duplicate: already committed, just ack it
charge_card(msg) # same transaction, so both commit or neither does
broker.ack(msg) # only after the commit — never before
Two details carry the weight. The dedupe row and the effect commit in one transaction, so there is no window where the system believes it processed something it did not. And the ack comes after the commit, which is exactly what makes this at-least-once rather than at-most-once. If charge_card is an external API rather than a database write, push the idempotency key to it — every serious payments API accepts one for precisely this reason.
Interviewers reliably probe this by asking "what if the consumer crashes here?" Having a specific answer that names the key and where uniqueness is enforced is worth more than any amount of broker trivia.
Ordering, partitions, and what breaks when you scale out
Global ordering and horizontal throughput are in direct conflict: a total order requires a single serialization point, and a single serialization point is a single machine's throughput. Every log resolves this the same way — order is guaranteed within a partition, and nowhere else.
Three consequences follow, and interviewers ask about all of them.
Order is per key, and you choose the key. Key by user_id and that user's events arrive in order — enough for "profile updated then deleted" to behave. Key by order_id and per-order state transitions are safe while different orders interleave freely. You almost never need global order; say which key you are ordering by and why.
Parallelism is capped by partitions. In a Kafka consumer group, a partition is assigned to at most one consumer. Twelve partitions means at most twelve useful consumers; the thirteenth sits idle. Size it from the arithmetic: if you need 60,000 messages/second and one consumer instance handles 5,000/s, you need at least 12 partitions, and you would provision maybe 24 so you can scale out later without re-partitioning.
Repartitioning breaks key affinity. Kafka's default partitioner is hash(key) % numPartitions. Change the partition count and almost every key moves — the exact failure mode that motivates consistent hashing, except here it is worse than a cache miss: user-42's new events land in a different partition than their old ones, so ordering across the change is lost. Kafka also refuses to reduce partition count at all. Over-provision partitions up front; it is much cheaper than the migration.
The hot-partition failure is the same shape as a hot cache key. If 30% of your traffic is one tenant, keying by tenant_id puts 30% of the load on one partition and one consumer, and no amount of scaling out helps. Fix it with a composite key (tenant_id:bucket) and accept that you have given up ordering within that tenant — which is a trade you should name out loud rather than discover in production.
Backpressure: the queue turns an outage into a delay
A queue does not create capacity. It converts a rate mismatch into a time debt, and the arithmetic is unforgiving.
steady state: 8,000 msg/s produced, consumers handle 10,000 msg/s
spike: 40,000 msg/s for 5 minutes
backlog built = (40,000 - 10,000) x 300s = 9,000,000 messages
drain rate = 10,000 - 8,000 = 2,000 msg/s
time to drain = 9,000,000 / 2,000 = 4,500s = 75 minutes
A five-minute spike produces 75 minutes of lag. At 1KB per message the backlog is 9GB, which is fine on disk — the storage was never the problem. The problem is that anything downstream is now over an hour stale, and if those messages are password-reset emails, you have an outage that no dashboard shows as red.
That is why consumer lag, not queue depth, is the metric to alert on, and why "we'll queue it" is a complete answer only when you also say what the drain looks like. If the work is latency-sensitive, you need autoscaling keyed on lag, or you need to shed load at the producer instead of accepting it — a queue in front of an under-provisioned consumer is a very expensive way to fail slowly. The same reasoning drives rate limiting at the edge: it is better to reject work you cannot do than to accept it and be late.
The dual-write problem, and the outbox
Here is a bug that a message queue introduces, that almost no interview prep material covers, and that will happen to you in production.
Your service commits an order to Postgres and then publishes OrderPlaced. Those are two systems with no shared transaction. Crash in between, and the order exists but nothing downstream will ever know. Publish first instead and you get the mirror-image bug: an event for an order that was never saved.
The fix is the transactional outbox. Write the order row and a row in an outbox table in the same transaction. A separate relay reads unsent outbox rows and publishes them, marking each sent. Either both rows committed or neither did, so the event can never disappear. In production the relay is usually change data capture — Debezium tailing the Postgres WAL — rather than a polling loop, but polling SELECT ... FROM outbox WHERE sent_at IS NULL LIMIT 100 every 200ms is perfectly adequate at moderate volume and far easier to run.
Mentioning the outbox unprompted when you introduce a queue into a design is one of the strongest signals available in an interview, because it shows you are thinking about the failure between two systems rather than the happy path through them.
Poison messages and dead letter queues
One message that always throws will be redelivered forever, and in a partitioned log it blocks its partition — head-of-line blocking that stalls every message behind it. Every broker's answer is a dead letter queue: after N attempts, move the message aside and continue.
Two things to say about DLQs. Retry with exponential backoff and jitter, not a tight loop, or your retries become the denial of service against a service that is merely degraded. And a DLQ with no alert and no redrive path is a data loss mechanism with extra steps — the useful version has a dashboard, an alarm on non-zero depth, and a documented way to fix the bug and replay the messages.
Picking a broker
| Model | Ordering | Retention | Where it shines | |
|---|---|---|---|---|
| Kafka | Partitioned log | Per partition | Time/size based, days by default | High throughput, replay, many independent consumers, stream processing |
| RabbitMQ | Queue with exchanges | Per queue, single consumer | Until acked | Complex routing, per-message priority and TTL, RPC-style work dispatch |
| SQS | Managed queue | None (standard) / per group (FIFO) | Up to 14 days | Zero operations, spiky load, AWS-native glue |
| Redis Streams | Log with consumer groups | Per stream | Manually trimmed | You already run Redis and need something modest |
Numbers worth carrying: a single Kafka partition comfortably sustains tens of MB/s, and a modest cluster does millions of messages a second. RabbitMQ handles tens of thousands per second per queue and starts to struggle when queues grow deep, because it is designed for messages to pass through rather than pile up. SQS standard is effectively unlimited in throughput but gives no ordering and explicitly warns about duplicates; SQS FIFO is capped at 300 API calls per second per queue (3,000 messages/s with batches of ten) unless you enable high-throughput mode.
If I had to choose without further context: SQS if you are on AWS and the work is fire-and-forget task dispatch, because the operational cost of Kafka is real and you should not pay it for a job queue. Kafka once the same event has more than one consumer, or once anyone asks to reprocess history — that is the point where a queue's delete-on-ack model starts costing you more than Kafka costs to run. RabbitMQ when routing is genuinely the hard part: topic exchanges, priorities and per-message TTL are things Kafka does not have and you would end up hand-rolling badly.
What message queues do not solve
They do not make the work faster. They move it. If your consumers cannot keep up, a queue converts a fast error into a slow success, which is often worse — users see a spinner that resolves hours later rather than an error they can act on.
They do not give you transactions across services. A published event is not a distributed commit. If the downstream step can fail permanently — the payment is declined after the order is confirmed — you need a compensating action and a state machine, not a queue. Queues move messages; sagas handle the rollback semantics, and only one of those is a broker feature.
They do not remove coupling, they change its shape. Consumers still depend on the event schema, and now there is no compiler and no synchronous error to tell you when you break it. A schema registry with compatibility rules is not optional at any real scale; the alternative is finding out from a consumer's exception log.
They do not help request-response work. If the caller needs the answer, a queue adds a hop and a correlation-id round trip to reconstruct what an RPC gave you for free.
They make debugging materially harder. The causal chain is now spread across services, brokers and retries. Budget for distributed tracing with the trace context propagated in message headers, or accept that "why did this user get two emails" becomes an afternoon.
The alternatives worth naming
Naming one of these, with a reason, is a strong signal — it shows you reached for a queue deliberately rather than reflexively.
A database table as a queue. SELECT ... FOR UPDATE SKIP LOCKED in Postgres gives you a correct work queue with no extra infrastructure, transactional enqueue for free (no outbox needed — it is the same database), and trivial inspection with SQL. It scales to thousands of jobs per second, which is more than most systems ever need. For a small team this is very often the right answer, and saying so demonstrates better judgement than defaulting to Kafka.
Synchronous call with retries and a circuit breaker. If the work must complete before you respond and the dependency is usually fast, a bounded retry with a timeout is simpler and easier to reason about than an async round trip.
A durable execution engine — Temporal, AWS Step Functions, Restate. For multi-step workflows with compensation, timers and human approval steps, these replace a pile of queues, state columns and retry logic with code that resumes where it stopped. When the design has more than about three sequential async steps that can each fail, this is usually the better tool.
Change data capture straight from the database. Debezium turns the WAL into a topic without the application publishing anything, which eliminates the dual write entirely. The cost is coupling consumers to your table schema.
Batch. If the requirement is "within an hour", a cron job over a table is less machinery, easier to backfill, and easier to reason about than a streaming pipeline.
Answering it in an interview
When you introduce a queue, say five things in about thirty seconds. Most candidates say only the first.
- Why the work is deferrable. "The user does not need the transcode to finish before we respond."
- Queue or log, and why. "It is a fact several services will want, so a log — and we may need to replay it."
- The partition key. "Keyed by
video_id, so per-video events stay ordered; ordering across videos does not matter." - The delivery guarantee and the idempotency key. "At-least-once, deduped on
event_idwith a unique index in the consumer's database." - What happens when consumers fall behind. "We alert on lag, autoscale consumers, and dead-letter after five attempts."
Then volunteer the failure before you are asked: "the risk is the dual write between our database and the broker, so I'd use an outbox table in the same transaction." Interviewers are scoring whether you can find the hole in your own design — and on this topic the hole is always either the dual write or the backlog, so you may as well name one first.
If you want to be pushed on those follow-ups out loud — "what happens when that consumer is down for an hour?" is the one that always comes next — you can run a full system design mock by voice on Whitepad.
Practice this out loud
Reading is the easy part. Sit across from a senior AI interviewer that talks, watches your whiteboard, and scores you like the real thing — your first mock is free.
Start a free mock →