Overview
Curated: · Written: · Reviewed:
Queues decouple in time, and everything follows from that
A message queue lets a producer hand off work without waiting for a consumer. What that buys is decoupling in three dimensions: time, since the consumer need not be available when the message is sent; load, since a burst is absorbed by the queue and drained at the consumer's rate; and failure, since a consumer outage delays work rather than losing it.
What it costs is that the operation is no longer complete when the caller returns. Everything difficult about queues follows from that single fact — the caller cannot report the outcome, errors surface somewhere else, and the system's state is briefly inconsistent by design.
So the first question about any queue is not which broker but: what does the user see between enqueue and completion? If the answer has not been designed, the queue has moved a problem rather than solved one.
Delivery guarantees
At-most-once: the message may be lost, never duplicated. Acceptable for a metric sample; rarely acceptable for work.
At-least-once: the message will be delivered, possibly more than once. This is what essentially every practical system provides, because the alternative requires losing messages.
A broker alone cannot guarantee exactly-once external side effects end to end across independent systems. A consumer can process a message and fail before acknowledging it, and an acknowledgement can itself be lost. Closing that gap requires a shared transaction boundary or application-level deduplication/idempotency, not a delivery label by itself. Systems advertising it provide something narrower: exactly-once processing within their own boundary, usually via transactional offsets or deduplication over a window. That is genuinely useful and is not the same as exactly-once when your consumer writes to an external database.
The practical consequence: consumers must be idempotent. Processing the same message twice must produce the same end state. This is not a nicety, it is the fundamental requirement, and a design that omits it is broken regardless of broker.
Idempotency is achieved with a natural key — a stable identifier from the producer — plus a conditional write: an upsert keyed on it, or a processed-messages table checked in the same transaction as the work. A key derived from a timestamp or a generated identifier makes every retry look new and defeats the whole mechanism.
Ordering
Ordering guarantees are narrower than people assume. Kafka guarantees order within a partition, not across a topic. FIFO queues guarantee order within a message group. Neither guarantees global ordering, and both achieve their guarantee by serialising the ordered unit — which caps parallelism at the number of partitions or groups.
That trade is the important part: ordering and parallelism are opposed. Ordering everything means processing everything sequentially.
So the design question is what actually needs ordering. Usually it is per-entity — events for one account must be ordered relative to each other, and events for different accounts need not be. Partitioning by that entity gives the ordering that matters while allowing parallelism across entities. Requiring global ordering is almost always over-specification, and it is expensive.
Better still, make consumers order-insensitive where possible: include a version or timestamp and ignore anything older than what has been applied. That tolerates out-of-order delivery entirely, which matters because a retry reorders messages on any competing-consumer or standard queue. It is not universal: SQS FIFO blocks the rest of a message group while a message is in flight, and a Kafka consumer that does not advance the partition offset will not process N+1 before N. Ordered partitions and FIFO groups preserve processing order under retry unless the consumer skips ahead, so the question to ask is which of the two a given broker and consumer give you.
Failure handling
Retries need exponential backoff with jitter. Immediate retries against a struggling dependency add load precisely when it is least able to take it, and synchronised retries from many consumers produce a thundering herd.
Distinguish retryable from non-retryable failures. A timeout or a 503 should be retried; a validation error or a 400 will fail identically forever, and retrying it wastes capacity and delays real work.
Dead-letter queues hold messages that have failed repeatedly, so a single poison message cannot block a queue indefinitely. The critical operational point: a dead-letter queue that nobody monitors is a silent data loss mechanism. It needs an alarm on depth, a process for inspecting messages, and a way to replay them after a fix.
Visibility timeouts hide an in-flight message from other consumers. If processing exceeds the timeout, the message becomes visible again and is processed twice concurrently — so the timeout must exceed realistic processing time, and long-running work should extend it explicitly rather than hoping.
Brokers differ in kind
Queue brokers (SQS, RabbitMQ) treat messages as work: consumed, acknowledged, deleted. Suited to task distribution, with routing and per-message acknowledgement.
Log-based systems (Kafka) treat messages as an append-only log retained independently of consumption. Consumers track a position and can replay from any point, and several independent consumer groups can read the same stream. That replayability is the defining difference and it enables things a queue cannot: reprocessing history after a bug, adding a new consumer that catches up from the beginning, rebuilding a derived store.
The cost is that consumer parallelism is bounded by partition count, ordering is per-partition, and offset management is the consumer's responsibility.
A database table is a legitimate queue at moderate throughput, using row locking that skips locked rows so workers claim different jobs. Its real advantage is transactional enqueue with your business data, which no external broker can give you without an outbox. Its limit is polling load and row contention at high rates.
The dual-write problem
A database write and a message publish are separate systems with no shared transaction, so either can succeed while the other fails. Publish first and a rollback leaves a message about something that never happened; commit first and a crash loses the message. There is no safe ordering — that is the key insight, not a bug to be careful about.
The outbox pattern resolves it: write the business change and a row in an outbox table in the same transaction, atomically. A separate process publishes unpublished rows and marks them sent. It can crash and restart, so nothing is lost; it can publish twice, which is at-least-once, which consumers already handle.
What to design before choosing a broker
The user-visible behaviour between enqueue and completion. How a failure is surfaced to whoever cares. Whether ordering is genuinely needed and at what granularity. What makes each consumer idempotent. What monitoring exists on queue depth, consumer lag and dead-letter depth. And what happens when the queue backs up — because it will, and shedding load deliberately is better than discovering the behaviour during an incident.
Worked example: sizing a queue from arrival and service rates
A checkout service publishes 1,200 orders/second at peak. Each consumer handles one order in 40 ms, so one consumer serves 25 orders/second.
Little's Law gives the concurrency the fleet has to sustain: L = lambda x W = 1,200/s x 0.040 s = 48 handlers in flight. That is the break-even fleet, and running at exactly break-even is the mistake — utilization 1.0 means a queue that never drains.
What the extra consumers buy is not lower steady-state latency. With one shared queue and capacity above arrivals, the backlog sits near zero and messages are picked up as they land. What they buy is recovery speed, so the column worth reading is the drain time after a 60-second consumer outage, which leaves 60 x 1,200 = 72,000 messages queued:
| Consumers | Throughput | Utilization | Headroom | Drain 72,000 |
|---|---|---|---|---|
| 48 | 1,200/s | 1.00 | 0/s | never |
| 55 | 1,375/s | 0.87 | 175/s | 411 s |
| 60 | 1,500/s | 0.80 | 300/s | 240 s |
| 75 | 1,875/s | 0.64 | 675/s | 107 s |
Each drain is backlog divided by headroom, because arrivals keep coming at 1,200/s while the backlog clears: 72,000 / 175 is 411 seconds, 72,000 / 300 is 240. The five consumers between 55 and 60 are worth almost three minutes of recovery, which is the difference between an incident nobody notices and one that pages.
Bounding the queue is what makes the failure mode a choice rather than an accident. A bound of 100,000 holds 100,000 / 1,200 = 83 seconds of peak arrivals, so the 60-second outage above never reaches it — 72,000 queued, 28,000 of headroom left. Push the outage to 90 seconds and the bound starts doing its job:
t+0.0s Consumers stop. Arrivals continue at 1,200/s.
t+83.3s Queue reaches its 100,000 bound (100,000 / 1,200).
t+83.3s Publishes now fail immediately rather than waiting: checkout returns 503.
t+90.0s Consumers recover. 108,000 were offered during the outage and 100,000 fit,
so 8,000 orders were rejected -- 7.4%.
t+423.3s 60 consumers drain 100,000 at 300/s of headroom. Queue empty.
Fail-fast rather than block-then-timeout is the deliberate half of that design, and the arithmetic shows why. Had publishes instead blocked with a 30-second timeout, none of them would have failed: consumers return 6.7 seconds after the bound is reached, waiters then clear at 300/s, so the longest wait is about 27 seconds and every blocked publish gets in. That sounds better until you notice the cost — a synchronous checkout request held open for up to 27 seconds, occupying a connection and a customer's patience. On a user-facing path, failing in milliseconds is the kinder answer; blocking belongs where the publisher is a batch job that can genuinely wait.
Unbounded, none of these lines exist: broker memory grows until it evicts, crashes, or takes the publisher down with it, and the failure is found at 3am in the broker rather than at the checkout that caused it. The 8,000 rejected orders are a real loss; they are a smaller and far more attributable one.
Two numbers to have ready for a follow-up. At-least-once delivery means the 40 ms handler must be idempotent: at 1,200/s, even a 0.1% redelivery rate is 72 duplicate orders a minute. And a dead-letter queue needs its own alarm — a DLQ nobody reads converts a loud failure into a silent one.
