Message queues
A durable buffer between producers and consumers that hands each message to one consumer at a time and redelivers it unless acknowledged.
Communication
Learn it
A producer has work for a consumer, but the two run at different speeds and fail independently. If the producer calls the consumer directly, a slow or down consumer stalls the producer or loses the work.
A queue sits between them: it stores messages durably and hands each one to one consumer at a time.
The core loop of a work queue:
- A consumer receives a message. The queue hides it from others for a visibility timeout (a lease) instead of deleting it.
- The consumer processes it, then acknowledges it, which deletes it.
- If no acknowledgement arrives before the timeout (crash, hang, slowness), the message becomes visible again and another consumer receives it.
Check
A consumer finishes processing a message and crashes before acknowledging it. What does the queue do?Other properties vary by system, and matter more than the brand:
- Ordering: most work queues don't preserve order across consumers. FIFO queues and partitioned logs give order per key, at the cost of parallelism. See Ordering.
- Dead-lettering: after N failed receives, a message moves aside so a poison message stops consuming capacity.
- Retention: work queues delete acknowledged messages; logs (Kafka-style) keep them and let each consumer track its own offset, so many consumers can read the same stream.
Getting a message into a queue atomically with a database change is its own problem; see Transactional outbox.
Check
One message crashes every consumer that receives it. Without a dead-letter queue, what happens?
Quick reference
The same ideas, condensed for revision.
How it goes wrong
- Timeout shorter than processing
- Slow messages become visible again while still being processed, so two consumers work on them concurrently.
- Poison message
- A message that always crashes its consumer is redelivered forever unless it is dead-lettered.
- Ack before work
- Acknowledging on receipt turns at-least-once into at-most-once: a crash after the ack loses the message.
- Lost on publish
- Publishing after a database commit, without an outbox, loses messages when the process dies in between.
Instead, consider
- Jobs table with SKIP LOCKED
- Throughput is modest and you want enqueueing to be transactional with your other writes.
- Partitioned log (Kafka, Kinesis)
- Many independent consumers need the same events, or you need replay and per-key ordering at high throughput.
- Direct RPC with retries
- The caller needs the result now and both sides are reliably up.
In practice
- In-process queue
- A bounded channel in memory: buffering without durability.
- Postgres / MySQL jobs table
- SELECT … FOR UPDATE SKIP LOCKED plus a lease column.
- SQS, Cloud Tasks, RabbitMQ
- Managed work queues with visibility timeouts or acks, and dead-letter queues.
- Kafka, Kinesis, Redis Streams
- Retained logs with consumer offsets and per-partition order.
It assumes
- Consumers can tolerate redelivery: processing is idempotent or deduplicated.
- The visibility timeout exceeds the normal processing time, or is extended by heartbeats.
- Something watches queue age and the dead-letter queue.
Explain it in your own words
Where you practise it
- A reliable video processing pipeline
How do jobs reach workers? · 100x traffic: find the real bottleneck · Paid instructors want a fast lane
- A job queue that keeps working when workers fall behind
What the numbers say · Why the queue stopped draining · A buffer that can hold a bad day · One slow job type · Twenty million jobs waiting · Defend keeping Redis
- Product analytics over billions of events
The database is down for twenty minutes · Defend the pipeline
- Notifications across email, push and in-app
Trace a mention to a phone · The announcement and the security alert
Further reading
Engineers describing it in systems they run.
- Scaling Slack's Job Queue
Slack · Saroj Yadav and others · Post, Dec 2017
An outage, the design flaw behind it, and the redesign. The section on rolling the new path out without an outage is worth reading twice.
Related concepts
- Delivery guarantees
At-most-once, at-least-once, and why 'exactly-once' is achieved by making duplicates harmless rather than by preventing them.
- Transactional outbox
Recording outgoing messages in the same database transaction as the state change, then delivering them separately, to avoid the dual-write problem.
- Leases and fencing tokens
Ownership that expires unless renewed, plus a token that lets the rest of the system reject an owner that has lost its claim without knowing it.
- Asynchronous processing
Accepting a request, recording the work durably, and doing it later in another process, so the work can outlive the request.
- Ordering
There is no global 'now' in a distributed system. Order exists only where something assigns it, so decide which order you need and who assigns it.
- Publish/subscribe
Decoupling senders from receivers by topic: a publisher sends once and every current subscriber receives a copy.
- Backpressure and capacity
When work arrives faster than it can be done, something has to give: the queue grows, the producer slows, or work is shed. Choose which on purpose.