Messaging & Queues

Decouple services using asynchronous communication patterns and message brokers.

Last generated

Lesson 6 of 18 available16 practice questions

SPACED REPETITION Β· 16 practice questions

Make this lesson stick.

Try 3 questions now. No account needed. Sample answers aren't saved.

A queue is a promise to finish the work later

Picture checkout at an online shop. The handler does six things in a row: authorize the card, reserve the stock, send the confirmation email, add loyalty points, record analytics and ask the warehouse to pick the order. Each call waits for the one before it.

On a normal day that works. Now run a flash sale. Assume 50,000 orders in the first 60 seconds, which is 833 orders per second. The email provider accepts 100 sends per second. The warehouse's legacy system accepts 500 pick requests per second.

Predict first: what breaks first, and what does the buyer see?

Check your answer

The email call. Against a limit of 100 per second, about 88% of the 833 email calls each second are rejected (1 βˆ’ 100/833). If a failed email fails the checkout, 88% of buyers get an error page after their card was authorized. If you swallow email errors instead, the warehouse breaks next: it can take 500 of the 833 requests per second, so 333 more pile up every second, threads block and timeouts spread back to the load balancer. Either way, a service that sells nothing decides whether checkout works.

The same chain costs you even on a quiet day:

  • Latency adds up. With assumed latencies of 350 ms (payment), 40 ms (stock), 800 ms (email), 60 ms (loyalty), 30 ms (analytics) and 120 ms (warehouse), the buyer waits 1.4 s. A full second of that is work they never see.
  • Availability multiplies. If each of the six services is up 99.9% of the time and they fail independently, checkout is up only 0.999⁢ β‰ˆ 99.4%. That is about 4.3 hours of failed checkouts a month instead of 43 minutes.

The fix is to split the work by one question: what must the buyer see before we answer? Whether the card was authorized, the stock reserved and the order created: yes. The email, the points, the analytics and the pick request must happen, but not before we answer. They go through a message broker. The checkout service hands the broker a message, and each downstream service processes it at its own pace.

BEFORE: all six calls on the request path
buyer ─► checkout ─► payment ─► stock ─► email ─► loyalty ─► analytics ─► warehouse
         (waits for all six; fails if any one fails)

AFTER: only what the buyer must see stays synchronous
buyer ─► checkout ─► payment (authorize) ─► stock (reserve)
            β”‚
            └─► database: order row + OrderPlaced event, one transaction   (~400 ms total)
                    β”‚
                    β–Ό
             [ message broker ] ─► email, loyalty, analytics, warehouse
                                    (each at its own pace; any may be down for a while)

A producer sends a message to a broker (Kafka, RabbitMQ, Amazon SQS, Google Pub/Sub…), which stores it until a consumer processes it. Producer and consumer never talk directly. So they don't need to be up at the same moment, run at the same speed, or know each other's addresses. That is the whole point, and it is also where every problem in this lesson comes from: once "later" is allowed, you must decide what happens when later means twice, out of order, an hour from now or never.

# Pressure in the requirements Move Signature example
1 The caller waits for work it doesn't need to see Take work off the request path Resize photos after upload
2 "Nothing may be lost", and machines crash Make duplicates harmless Loyalty points awarded once
3 Several teams react to the same event Fan out with pub/sub OrderPlaced β†’ email, loyalty, analytics
4 Per-entity order, high volume, or rereading history Keep a partitioned log Order status events keyed by order id
5 Arrivals outrun processing Absorb bursts, then push back Flash sale into a slow warehouse
6 Some messages will never succeed Retry with a limit, then dead-letter A malformed event

Scope: load balancers and the map of building blocks live in Core Building Blocks; database transactions in Databases & Storage; sagas in Advanced Topics & Final Prep. Here you learn the messaging decisions those designs rest on.

Before any move: what must the caller see, and when?

Before you draw a broker, ask five questions about each piece of work:

Ask Why it matters
Does the caller need the result to continue? Yes β†’ keep it synchronous. A queue cannot answer "did my payment go through?"
How late may it happen: a second, a minute, a day? Sets how much backlog you can tolerate and what you alert on
What if it happens twice? Decides how much idempotency work the consumer needs (Move 2)
Who else needs to know it happened? One consumer β†’ a queue; many β†’ pub/sub (Move 3)
Must it happen in order with related events? Decides the partition key (Move 4)

When a queue is the wrong answer. A page that shows the account balance needs the balance now. Routing that read through a queue adds a hop, a reply queue, correlation ids to match replies to requests, and a broker on the critical path, and it buys nothing. The same goes for a latency-critical hot path. A broker is also one more system to run, monitor and pay for. Reach for it when the work can be late, not by default.

When the client started something long. If a client uploads a video, don't hold the connection open for the transcode. Accept the job, return HTTP 202 Accepted with a status URL, and let the client poll or be notified. (Status codes are in Networking Basics.)

πŸ’‘ What to say: "Card authorization and the order write stay synchronous because the buyer must see the result. Everything else is a consequence of the order, so I publish an OrderPlaced event and let each downstream service consume it at its own pace." Likely follow-up: "What if the publish fails right after you commit the order?" Move 1 answers it.

Move 1: Take work off the request path

Pressure: the caller waits for work it doesn't need to watch, or the work is heavy, slow or bursty.

The task queue

A user uploads a profile photo. The web server stores the original in object storage, puts a small job message on a queue (ResizePhoto, with the photo id and the storage key) and answers 202 at once. A pool of workers takes jobs off the queue and resizes each photo into four sizes.

browser ─► web server ─► object storage (original photo)
               β”‚
               β”œβ”€β–Ί [ resize-jobs queue ] ◄─ worker 1
               β”‚                         ◄─ worker 2      each job goes to ONE worker at a time
               β”‚                         ◄─ worker 3
               └─► 202 Accepted + status URL

Several workers reading one queue are competing consumers: each message goes to one of them at a time. That is how a queue spreads work. Need more throughput? Add workers.

Ack after the work, not before

The broker hands a job out but forgets it only when the worker says the work is done: an acknowledgment (ack). Here is the trace with Amazon SQS, where a received message becomes invisible for a visibility timeout (30 s by default) and reappears unless the worker deletes it:

Time Worker A Worker B Message m in the queue
0 s receives m invisible until 30 s
4 s crashes mid-resize still invisible
30 s visible again
31 s receives m invisible until 61 s
33 s finishes, deletes m gone

Other brokers do the same thing in their own way. RabbitMQ returns an unacked message to the queue when the consumer's channel or connection closes. A classic Kafka consumer group keeps no per-message state at all: a restarted consumer resumes from the last offset it committed (Move 4). The idea is the same everywhere: only forget work that is finished.

Where you put the ack decides what a crash costs:

  • Ack before the work (RabbitMQ auto-ack, deleting a message as soon as it is received, committing a Kafka offset before processing): a crash loses the message. That is at-most-once.
  • Ack after the work: a crash after the work but before the ack runs the job again. That is at-least-once, and it is why Move 2 exists.

Predict: the visibility timeout is 30 s, and a large photo takes 45 s to resize. What happens?

Check your answer

At 30 s the message becomes visible again and a second worker picks it up while the first is still working. The resize runs twice and both workers write the results. The queue cannot tell a slow worker from a dead one. Set the timeout above your maximum processing time, not the p99, or have the worker extend it while it works (SQS ChangeMessageVisibility). Make the job idempotent anyway, because some duplicate will get through eventually.

Push or pull. RabbitMQ pushes messages to consumers but caps the number of unacknowledged messages each consumer holds with a prefetch count. SQS and Kafka consumers pull: SQS long polling waits up to 20 s for a message to arrive. Either way, a consumer should never take more than it can finish. Prefetch is backpressure in miniature (Move 5).

How many workers?

Little's law (Key Concepts & Terminology) says jobs in flight = arrival rate Γ— time per job. Assume 120 uploads per second at peak, 1.5 s of work per job, and 8 jobs in parallel per worker machine.

Step Arithmetic Result
Jobs in flight 120/s Γ— 1.5 s 180
Machines at 100% busy 180 Γ· 8 22.5 β†’ 23
Machines at a 70% utilization target 22.5 Γ· 0.7 32.1 β†’ 33

Your turn: the product team adds video thumbnails: 40 jobs per second, 6 s each, 4 at a time per machine. How many machines at a 70% target? And what does that mean for the visibility timeout?

Check your answer

40 Γ— 6 = 240 jobs in flight; 240 Γ· 4 = 60 machines fully busy; 60 Γ· 0.7 = 85.7, so 86 machines. The visibility timeout (or the heartbeat that extends it) must sit well above the slowest thumbnail job, not the 6 s average. Otherwise slow jobs run twice and quietly add load.

The hole between your database and the broker

Back to checkout. The handler commits the order row, then publishes OrderPlaced. If the process dies between those two steps, the order exists and nothing downstream ever hears about it: no email, no pick request. Swap the order (publish, then commit) and a rolled-back transaction leaves consumers acting on an order that doesn't exist. Two systems, no shared transaction: this is the dual-write problem.

The standard fix is the transactional outbox. Write the event into an outbox table in the same database transaction as the order. A separate relay reads unsent outbox rows, publishes them and marks them sent. The relay can poll the table or read the database's change log (change data capture, for example with Debezium).

checkout ─► BEGIN
              INSERT orders (...)
              INSERT outbox (event_id, type, payload)     ◄─ both or neither
            COMMIT

relay:  read unsent outbox rows ─► publish to broker ─► mark row sent
        (crash after publish, before mark β†’ publishes the row again)

The order and its event can no longer disagree. The price is a table, a relay to run, a small delay, and duplicates: a relay that crashes after publishing but before marking the row will publish it again. So the outbox gives you at-least-once publishing, and consumers must cope. Sagas and the outbox in depth: Advanced Topics & Final Prep.

⚠️ The card authorization is a side effect outside your database too. Give it an idempotency key derived from the order, so a retried checkout cannot charge twice. Move 2 shows why.

πŸ’‘ What to say: "Checkout writes the order and an outbox row in one transaction; a relay publishes the outbox. That gives at-least-once publishing without a distributed transaction, so every consumer must be idempotent."

Move 2: Make duplicates harmless

Pressure: "no order may be lost", and networks, workers and relays fail.

Three promises, defined by where the ack sits

Promise How you get it What a crash does Fits
At-most-once Ack or commit before the work, or send without retrying Loses the message Perishable data: a live location ping every 4 s, a sampled metric
At-least-once Ack after the work; producers retry until confirmed Repeats the message Business events: orders, payments, sign-ups
Exactly-once (of the effect) At-least-once + an idempotent consumer, or a transaction inside one system Nothing, within its boundary Anything where a duplicate costs money

At-least-once is the usual choice for business data, but you have to set it up: manual acks after processing and publisher confirms in RabbitMQ; in Kafka, keep the producer default acks=all (the default since Kafka 3.0; don't lower it to 0 or 1) and commit offsets only after processing. RabbitMQ with auto-ack is at-most-once. A Kafka consumer that auto-commits offsets (every 5 s by default) while another thread is still processing the records can commit work that never finished, and lose it.

What "exactly-once" really covers. Kafka's transactions let a read-process-write loop that goes from Kafka topics to Kafka topics commit its output and its input offsets atomically. SQS FIFO queues drop a resend that carries the same deduplication id within 5 minutes. Kafka's idempotent producer, on by default since Kafka 3.0, stops the producer's internal retries from writing duplicates into the log, but not a fresh send from a restarted process. None of these reach outside their own system. An email, a card charge or a row in your PostgreSQL database is outside. The cost is not the problem: Confluent measured the transactional producer at about 3% below in-order at-least-once (1 KB messages, 100 ms transactions) (source). Scope is.

Where duplicates come from, in a system built only from correct parts:

  • A producer retries after its ack was lost; the broker already stored the first copy.
  • An outbox relay crashes between publishing and marking the row sent.
  • A consumer crashes after the side effect but before the ack.
  • A visibility timeout expires while the job is still running.
  • A Kafka consumer group rebalances before offsets were committed.
  • Someone replays a topic on purpose.

Duplicates are not bad luck. They are the normal price of not losing messages. The consumer has to make them harmless.

The idempotent consumer

Predict: a loyalty consumer checks a processed table for the event id, awards the points, then inserts the id. Find two ways it can award the points twice.

Check your answer
  1. Two deliveries at once. A visibility timeout expires or a rebalance hands the message to a second consumer while the first is still working. Both checks run before either insert, so both award the points.
  2. A crash between the effect and the insert. The points are committed, the process dies, the message is redelivered, the check finds nothing, and the points are awarded again.

Check-then-act is two steps, and a failure or a race can land between them. Recording the id first in its own transaction doesn't fix it either. It swaps duplicates for losses: crash after the insert and the redelivery is skipped, so the points never arrive.

The fix is to make the claim and the effect one atomic step: insert the event id and apply the change in the same database transaction. Here it is in runnable Python, with SQLite standing in for your database:

import sqlite3

db = sqlite3.connect(":memory:", isolation_level=None)   # we write BEGIN/COMMIT ourselves
db.executescript("""
    CREATE TABLE points    (customer_id TEXT PRIMARY KEY, balance INTEGER NOT NULL);
    CREATE TABLE processed (event_id    TEXT PRIMARY KEY);
    INSERT INTO points VALUES ('c-7', 0);
""")

def handle_order_paid(event, crash_before_commit=False):
    """Award loyalty points once per OrderPaid event, however often it arrives."""
    db.execute("BEGIN IMMEDIATE")
    try:
        claimed = db.execute("INSERT OR IGNORE INTO processed (event_id) VALUES (?)",
                             (event["event_id"],)).rowcount
        if claimed == 0:                      # already applied by an earlier delivery
            db.execute("ROLLBACK")
            return "duplicate, skipped"
        db.execute("UPDATE points SET balance = balance + ? WHERE customer_id = ?",
                   (event["points"], event["customer_id"]))
        if crash_before_commit:
            raise RuntimeError("worker killed")   # simulate a crash mid-handler
        db.execute("COMMIT")                  # the claim and the effect land together
        return "applied"
    except Exception:
        db.execute("ROLLBACK")                # ...or neither of them does
        raise

event = {"event_id": "order-981/paid", "customer_id": "c-7", "points": 40}
try:
    handle_order_paid(event, crash_before_commit=True)   # delivery 1 dies
except RuntimeError as e:
    print("delivery 1:", e)
print("delivery 2:", handle_order_paid(event))           # broker redelivers
print("delivery 3:", handle_order_paid(event))           # a duplicate arrives
print("balance:", db.execute("SELECT balance FROM points").fetchone()[0])
Delivery processed before What happens Balance after
1 empty claims the id, adds 40, crashes; both writes roll back 0
2 empty claims the id, adds 40, commits 40
3 order-981/paid claim finds the id: skip, then ack 40

The script prints worker killed, applied, duplicate, skipped and balance: 40.

Why it works. The primary key makes the claim atomic. In PostgreSQL you write INSERT … ON CONFLICT DO NOTHING, and a second transaction inserting the same key waits for the first one, then finds the conflict. Because the claim and the effect commit together, there is no moment when one exists without the other. A crash before the commit rolls both back, so the redelivery applies the change. A crash after the commit but before the ack leaves the claim in place, so the redelivery skips and acks.

Choose the key from the business operation, not the delivery. order-981/paid names one real-world fact. A random UUID generated each time the producer calls send does not: if the producer retries after a lost ack, the retry carries a new id and the check lets it through. Create the event id once, for example in the outbox row, and reuse it on every retry.

Keep keys as long as a duplicate can still arrive. A duplicate can arrive as late as the longest redelivery or replay window you allow. If you may replay 14 days of a topic, 24 hours of keys is not enough. The store is usually cheap: 3 million events a day Γ— 14 days Γ— about 100 bytes per key is roughly 4.2 GB.

Side effects you can't roll back

Card charges, emails, SMS messages and calls to partner APIs cannot join your database transaction.

  • Pass the key downstream when the provider supports it. Stripe's API takes an Idempotency-Key header, saves the result of the first request with that key, and returns the same result for retries. Stripe may prune keys once they are 24 hours old, so this covers retries, not a replay weeks later (Stripe docs).
  • Otherwise record the intent. Claim the key with status pending, make the call, then set done. A redelivery that finds pending cannot know whether the call happened, so it asks the provider (look up your reference) or accepts a rare duplicate.
  • Decide which duplicates you can live with. A second "your order shipped" email is annoying. A second charge is a support ticket and a refund. Say which is which.

Naturally idempotent operations need no bookkeeping. SET status = 'shipped' can run twice safely; SET stock = stock - 1 cannot. An upsert by key is safe. So is "apply only if event.version is newer than the stored version" when each event carries the full new state, and it also throws away stale events that arrive out of order.

πŸ’‘ What to say: "Delivery is at-least-once, so the loyalty consumer claims the event id in the same transaction as the points update. For the card charge I pass an idempotency key derived from the order to the payment provider." Likely follow-up: "And the email? You can't roll that back."

Your turn: the loyalty handler must also email "you earned 40 points". Where does the send go?

Check your answer

Not inside the transaction. If the email is sent and the commit then fails, the customer is told about points they don't have, and the retry sends a second email. Instead, write a PointsAwarded outbox row in the same transaction as the points. A separate email consumer sends from that event, keyed by its event id, with the provider's deduplication if it has one. The worst case left is a rare duplicate email, which is acceptable for this message. A duplicate charge would not be.

Move 3: Fan out with pub/sub

Pressure: several teams must react to the same event, and the producer shouldn't have to know who they are.

The loyalty team's missing points

Predict: email, loyalty and analytics all read from one SQS queue, order-events. The loyalty team reports that about two-thirds of customers never got their points. What happened?

Check your answer

They are competing consumers on one queue, so each message goes to only one of the three services. Loyalty sees roughly a third of the orders, depending on how fast each service polls, and so do the others: every service is missing orders, loyalty just noticed first. A queue is point-to-point. Several independent readers need a topic.

There are two delivery models, and mixing them up is a classic interview slip:

  • Point-to-point (a queue): each message goes to one consumer. Use it for commands, instructions with one handler: "resize photo 42", "send receipt 981".
  • Publish/subscribe (a topic): every subscriber gets its own copy. Use it for events, facts about the past that any number of services may care about: "order 981 was placed".

In a real design you use both, at two levels. You fan out across services, and the instances within each service compete:

                          β”Œβ”€β–Ί [ queue: email ]     ─► email workers Γ—3
OrderPlaced ─► [ topic ] ─┼─► [ queue: loyalty ]   ─► loyalty workers Γ—2
                          └─► [ queue: analytics ] ─► analytics workers Γ—4

      one copy per subscription; the workers inside one service compete
Broker Fan-out across services Competing workers within a service
AWS SNS topic β†’ one SQS queue per service several workers poll that service's queue
RabbitMQ fanout or topic exchange β†’ one bound queue per service several consumers on that queue
Kafka one consumer group per service, all reading the topic the group's instances split the partitions
Google Pub/Sub one subscription per service several subscribers pull that subscription

Giving each service its own queue or subscription also gives it its own backlog, retry policy and dead-letter queue. A slow analytics deploy no longer delays anyone's receipt email.

Adding a subscriber means adding a subscription. The producer doesn't change. But with an SNS standard topic and SQS, or with RabbitMQ queues, a new queue only receives messages published after it was created. (SNS FIFO topics can archive messages and replay them to a new subscription, if you enable and pay for that.) If the newcomer needs last month's orders, you need a log that keeps history (Move 4), an archive, or a backfill from the database.

Filtering. SNS subscriptions can filter on message attributes. RabbitMQ topic exchanges route by pattern: a queue bound to order.*.eu gets EU order events only. Kafka consumers read the whole topic and skip what they don't need, or you split the stream into separate topics.

What fan-out costs. Every subscriber is another copy: 2,000 events per second with 12 subscribers is 24,000 deliveries per second. The producer can no longer see who depends on its events, so its message format becomes a public contract (the contract section below). And a chain of events is harder to debug than a call stack, so carry a correlation id through every message. The same fan-out idea at feed scale, with celebrity accounts, is in Twitter/Instagram News Feed.

πŸ’‘ What to say: "OrderPlaced is an event, so it goes to a topic. Each consuming service gets its own queue with its own DLQ, and workers compete inside it. Adding a consumer doesn't touch checkout." Likely follow-up: "A new fraud team needs last month's orders. How?"

Your turn: analytics wants every order. The EU fraud team wants only EU orders above 1,000 EUR. How do you wire both without changing checkout?

Check your answer

Two subscriptions on the same topic. Analytics subscribes with no filter. Fraud subscribes with a filter: an SNS filter policy on region and amount attributes that the producer sets on each message, or a RabbitMQ binding on a routing key such as order.placed.eu with the amount check in the consumer. The cost of filtering at the broker is that the producer must publish the attributes you filter on, which makes them part of the contract.

Move 4: Keep a partitioned log

Pressure: events for the same entity must be applied in order, volume is high, or someone needs to reread history.

Queue versus log

A queue forgets a message once it is acked. A log appends records and keeps them for a retention period. Each consumer only remembers its position, the offset, and reading does not remove anything. Kafka is the best-known log. Amazon Kinesis and RabbitMQ Streams work the same way.

topic "order-status", 3 partitions, key = order_id

partition 0: [0][1][2][3][4][5][6]    billing group: instance A, next offset 5
partition 1: [0][1][2][3][4]          billing group: instance A, next offset 4
partition 2: [0][1][2][3][4][5]       billing group: instance B, next offset 2

analytics group: one instance reads all three partitions at its own offsets

Five rules cover most interview questions about Kafka:

  1. Same key, same partition. The default partitioner hashes the record key to pick the partition, so all events for order-981 land in one partition, in the order they were written. There is no order across partitions.
  2. A partition has one reader per group. In a classic consumer group, each partition is assigned to exactly one instance. Partitions therefore cap parallelism: in a 20-instance group reading 12 partitions, 8 instances sit idle.
  3. Each group keeps its own offsets. Billing and analytics read the same records independently. That is how Kafka does pub/sub.
  4. Commit after processing. A consumer commits its offset once the records are processed, so a crash or a rebalance restarts from the last commit and repeats everything after it. That is at-least-once again. A consumer that takes longer than max.poll.interval.ms (5 minutes by default) between polls is removed from the group, which triggers a rebalance.
  5. Retention is a clock, not an ack. Records stay for the retention period (7 days by default, log.retention.hours=168) or until a size limit, whether or not anyone read them. A compacted topic instead keeps at least the latest record for each key. A new consumer group starts at the end of the log by default (auto.offset.reset=latest). To replay history, start it at earliest or seek to a timestamp.

Pick the key from what must stay in order: order_id for order status, account_id for a ledger. A key that is too coarse creates hot partitions. No key at all spreads records across partitions with no per-entity order.

Why not add partitions later when you need them? Because the key-to-partition mapping moves:

import zlib

def partition_for(key: str, partitions: int) -> int:
    # Kafka's default partitioner uses murmur2; crc32 is a stand-in stable hash
    return zlib.crc32(key.encode()) % partitions

keys = [f"order-{i}" for i in range(100_000)]
moved = sum(partition_for(k, 12) != partition_for(k, 16) for k in keys)
print(f"{moved / len(keys):.0%} of keys change partition going from 12 to 16")
print(partition_for("order-1042", 12), partition_for("order-1042", 12))  # same key, same partition

It prints 75% of keys change partition going from 12 to 16. For three keys in four, the old events sit in one partition and the new events go to another, so a consumer can read a key's new events before it has finished the old ones. Choose a partition count with headroom up front.

Estimate: partitions and disk

Assume order-status events average 20,000 per second with peaks of 60,000, 1 KB each, and that one consumer instance can apply 2,500 events per second (it is limited by its database writes).

Quantity Arithmetic Result
Consumer instances at peak 60,000 Γ· 2,500 24
Partitions at least 24; choose 48 so you can double without remapping keys 48
Data per day 20,000/s Γ— 1 KB Γ— 86,400 s 1.73 TB
7-day retention Γ— 7 12.1 TB
With replication factor 3 Γ— 3 36.3 TB

The classic slips: using the peak rate for storage (3Γ— too much), forgetting replication (3Γ— too little), and mixing up bits and bytes (8Γ— off).

Hot keys: the merchant who ate a partition

Say the order-status topic is keyed by merchant_id, and one merchant produces 30% of all events. At peak that is 18,000 events per second into a single partition (plus its share of everyone else's traffic), and only one consumer can read it, at 2,500 per second. That partition's lag grows by more than 15,500 events every second. The other partition on the same consumer starves too, while the rest share the remaining 42,000 per second with room to spare. Your options:

  • Key by something finer if that is all ordering needs. Per-order ordering is usually enough, so key by order_id.
  • Split the hot key into sub-keys (merchant-17#0 … merchant-17#7) and give up ordering across the whole merchant.
  • Isolate the tenant with its own topic and consumers.

Head-of-line blocking. A partition is processed in order, so one slow or failing record holds up every record behind it, including records for other keys. With the usual single-threaded consumer loop it is worse: while the handler is stuck, that instance processes none of its partitions. A queue doesn't have this problem, because other workers simply take other messages. Move 6 deals with the failing case.

Queue or log?

Queue (RabbitMQ queues, Amazon SQS) Log (Kafka, Kinesis, RabbitMQ Streams)
After a consumer acks the message is deleted the record stays until retention or compaction
Unit of work sharing the message: any idle worker takes the next one the partition: at most one reader per classic group
Fan-out one queue per subscriber one consumer group per subscriber
Ordering weak once many workers compete; SQS FIFO orders within a message group per partition, so per key
Replay no: acked means gone yes, within retention: reset the offsets
Per-message retry and DLQ built in (visibility timeout, nack, redrive policy, dead-letter exchange) you build it (retry and DLQ topics)
One slow message delays only the worker holding it delays its partition, and with a single-threaded consumer every partition that instance owns
Good for jobs and commands with uneven, per-message work event streams, many readers, high volume, replay

Products blur this line, so choose the semantics first and name a product second. RabbitMQ Streams (since RabbitMQ 3.9) are append-only logs that consumers can reread from an offset or a timestamp (docs). Kafka's share groups add queue-style per-record acknowledgment and delivery counting on Kafka topics, with more consumers than partitions. They were declared production-ready in Kafka 4.2 (release notes).

Ordering on SQS. A FIFO queue (its name ends in .fifo) delivers messages with the same MessageGroupId in order, one at a time, and it needs a deduplication id or content-based deduplication. FIFO throughput is limited: without batching and in the default mode, 300 transactions per second per API action. High-throughput mode raises that a lot (SQS quotas).

⚠️ MessageGroupId on a standard queue does not order anything. It turns on fair queues, which stop one noisy tenant from delaying the others, and messages in the same group may still be processed in parallel.

πŸ’‘ What to say: "I'll key the order-status topic by order id, so each order's events stay in order and different orders run in parallel. Peak is 60k per second and one consumer applies 2.5k, so I need 24 consumers; I'll create 48 partitions for headroom. Seven days of retention lets a new team replay recent history; if a full-week backfill is a requirement, I'll set 10 to 14 days for margin." Likely follow-up: "What if one merchant is 30% of your traffic?"

Your turn: during a backlog, the on-call engineer scales the billing group from 12 to 20 instances. The topic has 12 partitions. What happens, and what would actually add throughput?

Check your answer

Eight instances sit idle: each of the 12 partitions already has its one reader, and the scale-up also triggered a rebalance that paused consumption briefly. Throughput per partition only rises if each instance processes faster: batch the database writes, or process different keys in parallel inside one instance while keeping each key's order. Longer term, add partitions and accept the key remapping at a quiet moment, or plan the partition count for peak from the start.

Move 5: Absorb bursts, then push back

Pressure: arrivals outrun processing, either briefly (a spike) or for good (growth, or a slow dependency).

Buffering a burst

Back to the flash sale: 50,000 orders in 60 s, so 833 per second, feeding the warehouse consumer, which handles 500 per second. With a queue in between, nothing is rejected. Orders wait instead.

Time Arrived Processed Backlog An order arriving now waits
0 s 0 0 0 0 s
15 s 12,500 7,500 5,000 10 s
30 s 25,000 15,000 10,000 20 s
60 s 50,000 30,000 20,000 40 s
80 s 50,000 40,000 10,000 β€”
100 s 50,000 50,000 0 β€”

The backlog grows at 833 βˆ’ 500 = 333 per second, peaks at 20,000 when the burst ends, and drains at 500 per second, so it is empty at 100 s. The last order of the burst waits 40 s for its pick request. The queue has turned failed requests into waiting time.

Predict: outside sales, orders arrive at 50 per second, and someone sizes the warehouse consumer at exactly 50 per second "because that's the load". What happens after the next spike?

Check your answer

The backlog never drains. Drain time = backlog Γ· (service rate βˆ’ arrival rate), and here the difference is zero. Worse, any ordinary wobble in traffic adds to the backlog and nothing takes it away. A queue only absorbs spikes if consumers have headroom above the average arrival rate. The headroom decides how fast you recover.

When the buffer isn't enough

Now the overload is sustained: 1,200 messages per second in, 890 out. The backlog grows by 310 per second, which is 1.1 million an hour. A queue does not fix this. It only decides how it fails:

  • SQS deletes messages older than the retention period (4 days by default, 14 at most). Old work vanishes silently.
  • Kafka deletes old log segments on its retention clock whether or not a lagging group has read them. That group loses the records.
  • RabbitMQ raises a memory or disk alarm and blocks publishing connections, while consumers keep going (docs). Now the pressure reaches your producers, and their requests time out.
  • A bounded queue (for example a RabbitMQ max-length) rejects new messages or drops the oldest, whichever you configured.

Backpressure is the mechanism that lets a slow stage slow down the faster stage in front of it. It is not the overload itself. You have four levers, each with a price:

Lever How You give up
Scale consumers autoscale on backlog or lag money; capped by partition count and by the database behind the consumers
Push back bound the queue; producers get 429 or 503 with Retry-After, or block callers see errors or latency
Shed load drop or sample low-value work (analytics events, not orders) completeness of that data
Prioritize separate queues for critical and bulk work more queues to run and watch

Rejecting with 429 is rate limiting; token buckets and friends are in Design URL Shortener & Rate Limiter.

Measure time, not just counts

A backlog of 500,000 is 50 seconds of work at 10,000 per second, or two hours at 70 per second. Count alone says nothing, so turn it into time:

  • A message arriving now waits about backlog Γ· drain rate.
  • Catching up takes backlog Γ· (drain rate βˆ’ arrival rate).
  • Watch the age of the oldest message directly where the broker exposes it (SQS publishes ApproximateAgeOfOldestMessage), and per-partition lag in Kafka.
queue: order-events               depth 12,450  (rising)
in:  1,200 msg/s                  out: 890 msg/s  (4 consumers)
oldest message age: 10 s          DLQ depth: 0

Predict: every number here looks harmless. The requirement is that an order event is processed within 2 minutes. How long until that breaks?

Check your answer

About 5 minutes. A new arrival waits depth Γ· drain rate = 12,450 Γ· 890 β‰ˆ 14 s today (the oldest message, at the head, has waited about depth Γ· arrival rate β‰ˆ 10 s). The wait reaches 120 s when the depth reaches 120 Γ— 890 = 106,800. At +310 per second from 12,450, that takes (106,800 βˆ’ 12,450) Γ· 310 β‰ˆ 304 s. The alarm on this dashboard is the growth rate, not the depth.

Autoscale on backlog or age, not on consumer CPU. A consumer blocked on a slow database shows low CPU while its backlog explodes. For SQS workers, AWS's guidance is a custom "backlog per instance" metric. On Kubernetes, KEDA scales pods on queue length or Kafka lag. Two ceilings stay: a classic Kafka consumer group will not use more consumers than partitions, and the database behind the consumers has its own limit.

Your turn: drain that backlog once it has reached 1.1 million, within 10 minutes, while 1,200 per second keep arriving. Each consumer does 890 Γ· 4 = 222.5 per second. How many consumers, and what else must you check?

Check your answer

Required rate = arrivals + backlog Γ· time = 1,200 + 1,100,000 Γ· 600 β‰ˆ 3,033 per second. 3,033 Γ· 222.5 = 13.6, so 14 consumers. Then check the ceilings: whatever the consumers write to must take about 3,000 writes per second, and on Kafka the topic needs at least 14 partitions.

Move 6: Retry with a limit, then dead-letter

Pressure: some messages fail. Some will succeed later: a timeout, or a 503 from a payment provider during an incident. Some never will: a missing field, a reference to a deleted customer, a bug.

A poison message fails every time. Without a limit it loops forever, burning workers, and on an ordered partition it blocks everything behind it.

Predict: a RabbitMQ consumer nacks with requeue=true on any exception, and a malformed message arrives. What do the graphs show?

Check your answer

A hot loop. A requeued message comes back immediately, fails again and is requeued again, thousands of times per second: CPU pegged, logs flooded, and the messages behind it wait. RabbitMQ's quorum queues put a ceiling on this: since RabbitMQ 4.0 their default delivery limit is 20, after which the message is dropped, or dead-lettered if you configured a dead-letter exchange (docs). Relying on that default still means 20 useless attempts, and possibly silent loss.

Four rules:

  1. Classify the failure. Transient (timeouts, 429, 503) β†’ retry. Permanent (validation error, unknown schema) β†’ dead-letter now; retrying cannot fix it.
  2. Delay retries with capped exponential backoff and jitter. Thousands of messages that failed together should not retry together.
  3. Cap attempts or total time, then dead-letter.
  4. Treat the dead-letter queue (DLQ) as an incident queue. Keep the failure reason and the original metadata, alert when it is non-empty for critical flows, fix the cause, then redrive the messages.
import random

def retry_delay(attempt: int, base: float = 1.0, cap: float = 300.0) -> float:
    """Seconds to wait before retry number `attempt` (1, 2, 3, ...):
    capped exponential backoff with full jitter."""
    ceiling = min(cap, base * 2 ** (attempt - 1))
    return random.uniform(0, ceiling)

for attempt in (1, 2, 3, 6, 10, 12):
    print(attempt, round(min(300.0, 2 ** (attempt - 1)), 1), round(retry_delay(attempt), 1))

The ceiling doubles from 1 s (1, 2, 4, 8 … 256 s) and then stays at 300 s. Over 10 retries the ceilings add up to 811 s, so a message waits at most about 13.5 minutes in total (about half that on average, because of the jitter) before it is dead-lettered. Full jitter picks a random delay between zero and the ceiling. That spreads out the retries of messages that failed at the same moment, so a dependency that has just recovered is not flattened by a synchronized wave: the thundering herd.

main queue ─► consumer ──ok──────────────────────────► ack
                 β”‚
                 β”œβ”€ transient ─► retry later: 1 s, 2 s, 4 s … cap 300 s, jittered
                 β”‚                  └─ attempts used up ─┐
                 └─ permanent ──────────────────────────►│
                                                          β–Ό
                                   dead-letter queue ─► alert, inspect, fix, redrive

How each broker delays a retry:

  • SQS: a message reappears after its visibility timeout, and you can change that per message (ChangeMessageVisibility) to back off. A redrive policy moves it to a DLQ after maxReceiveCount receives.
  • RabbitMQ: nack with requeue is immediate. For backoff, publish to a delay queue whose message TTL expires into a dead-letter exchange that routes the message back. A dead-letter exchange also catches rejected messages.
  • Kafka: a classic consumer group has no per-record retry. A record is read again only after a crash or rebalance, from the last committed offset. So the consumer retries in-process, which blocks its partition (and, in a single-threaded consumer, every partition that instance owns), or it publishes the record to a retry topic (for example orders.retry.1m) or a DLQ topic, then commits and moves on. Share groups (Move 4) count deliveries and redeliver per record instead.
  • Google Pub/Sub: subscriptions have a retry policy with exponential backoff and dead-letter topics with a maximum number of delivery attempts.

Stop calling what is down. Every attempt against an endpoint that times out after 10 s holds a worker for 10 s, backoff or not. When one downstream (a partner, a merchant's webhook, a provider) keeps failing, open a circuit breaker for it: for a while, reschedule its messages without calling it, then let a single probe test whether it has recovered. That keeps one dead dependency from eating the whole worker pool.

Retries multiply. If each of three layers retries 3 times, one user request can become (1 + 3)Β³ = 64 calls to the bottom service. That happens exactly when it is already struggling. Retry at one layer, or give each caller a retry budget.

⚠️ Dead-lettering breaks order. Say OrderShipped for order 77 goes to the DLQ, then OrderDelivered for order 77 arrives and is applied. The order is now "delivered" but never "shipped". For ordered streams, either park the key (hold later events for order 77 until the parked one is resolved), or give each key's events a sequence number and have the consumer hold any event that arrives after a gap. A plain "newer version wins" check does not help here: OrderDelivered is newer.

πŸ’‘ What to say: "Transient errors retry with capped exponential backoff and jitter for up to about 15 minutes; validation errors go straight to a DLQ. A non-empty DLQ pages the owning team, and we redrive after the fix." Likely follow-up: "What if the dead-lettered message was step 2 of an ordered sequence?"

Messages outlive deploys: the contract

A function call and its caller change in one deploy. A message may be read hours later, or replayed weeks later, by consumers running other versions of the code. So the message format is a public API. Treat it like one.

A useful envelope:

{
  "event_id": "evt_01J9ZK3Q8X",
  "event_type": "OrderPlaced",
  "schema_version": 2,
  "occurred_at": "2026-09-24T10:31:00Z",
  "producer": "checkout",
  "correlation_id": "req_7f3a",
  "key": "order-981",
  "data": { "order_id": "order-981", "customer_id": "c-7", "total_minor": 14999, "currency": "EUR" }
}
  • event_id is created once (in the outbox row) and reused on every retry. It is the consumer's dedupe key.
  • key decides the partition, and so the ordering (Move 4).
  • correlation_id threads one user request through every hop, so you can trace it.
  • schema_version changes only for breaking changes. Adding an optional field must not need a new code path.

Compatibility has two directions. Backward: new readers can read old messages, which matters when you replay history. Forward: old readers can read new messages, which matters during a rolling deploy where producers upgrade first. Adding an optional field with a default keeps both, if readers ignore fields they don't know. A strict deserializer that rejects unknown fields turns every addition into an outage. Renaming a field, removing it, or changing its type or meaning breaks readers.

Rename with expand/contract. Add the new field, write both, move every consumer to the new one, stop writing the old one, and remove it once no consumer reads it and no retained message has only the old field.

Formats. JSON is readable and everywhere, but larger, and its schema exists only by convention (or JSON Schema). Avro and Protobuf are compact, binary and schema-driven. A schema registry (Confluent Schema Registry supports Avro, Protobuf and JSON Schema) can reject an incompatible schema when the producer registers it, before any consumer breaks.

Size. Keep messages small. SQS caps a message at 1 MiB, and Kafka's default maximum is about 1 MB. For a video or a large document, store the object in object storage and send a pointer: the claim check pattern.

Fat or thin events. A fat event carries the data consumers need (total, items), so they never call back to the producer, but more of your schema leaks. A thin event says "order 981 changed", and every consumer calls the order service to fetch the current state, which puts load back on the producer. Say which one you chose and why.

Expiry. Put a time-to-live only on perishable messages, such as a one-time password or a location ping. A payment-completed event is never too old to matter. Letting it expire because consumers were down for an hour means a paid order that is never fulfilled.

Final round: no label on the problem

Real prompts don't say which move they want. Start from what the caller must see, then find the pressure.

Challenge 1: the campaign sender

A marketing tool sends a campaign to 2 million recipients through an email provider that accepts 500 sends per second and returns 429 above that. A user can cancel a running campaign. Nobody may receive the same campaign email twice.

  1. How long does a full send take?
  2. What queue shape, and how do you respect the provider's limit?
  3. Delivery is at-least-once. How do you meet "never twice"?
  4. How does cancel work?
  5. What belongs in the DLQ?
Hint

Question 3 is the twist. Look again at the two wrong orders of "record the key" and "do the work" in Move 2, and ask which failure this product prefers.

Check your answer
  1. 2,000,000 Γ· 500 = 4,000 s, about 67 minutes. Say so, because it sets the user's expectations.
  2. One job per recipient (or per batch of 100) on a queue, a pool of sender workers, and a shared rate limiter at 500 per second (a token bucket). A 429 means back off. More workers will not help, because the provider is the bottleneck.
  3. Claim campaign_id + recipient_id with status sending before calling the provider, and pass the same key to the provider if it supports idempotency. On redelivery, a sending claim is skipped, not retried. That makes each email at-most-once. It is the right promise here, because a missed marketing email is cheaper than a duplicate, and at worst you lose the handful of sends in flight when a worker died. Say that trade-off out loud.
  4. A cancelled flag on the campaign, checked before each send. Queued jobs drain as no-ops, so you don't need to find and delete them.
  5. Permanent failures: invalid or hard-bounced addresses. Not 429s, which are the backpressure signal and should be retried with backoff.

Challenge 2: the search index that must keep up

A product catalog database takes 3,000 updates per second at peak. Search must show each product's latest version within 5 seconds. Twice a year the index mapping changes and the whole index must be rebuilt from scratch. Two updates to the same product can arrive milliseconds apart.

Check your answer
  • Capture: change data capture on the catalog (or an outbox) publishes every update to a log. A queue would not do: the rebuild needs history.
  • Key: product_id, so each product's updates stay in order.
  • Retention: a compacted topic keeps the latest record for every product, as long as every product was written to it at least once, so seed it with the CDC connector's initial snapshot (Debezium takes one by default). A rebuild then replays the topic from the beginning into a new index and switches the search alias, without rereading the whole database.
  • Consumer: upsert by product_id, only if the event's version is newer than the stored one. That is idempotent (Move 2) and ignores stale events.
  • Size: if one indexer applies about 500 updates per second in bulk requests, peak needs 3,000 Γ· 500 = 6 consumers. Create 12 or more partitions for headroom and faster rebuilds.
  • Alert: on the age of the oldest unindexed update in each partition (lag expressed as time), well before the 5 s target, not on the raw lag count.

Cheat sheet: pressure β†’ move β†’ cost

When the requirements say… Reach for You pay with
the caller doesn't need to wait a queue + workers, 202 Accepted (Move 1) eventual completion, a broker to run
an event and a DB write must agree a transactional outbox (Move 1) a relay, a small delay, duplicates
nothing may be lost ack after work, at-least-once (Move 2) duplicates
duplicates cost money an idempotent consumer: claim + effect in one transaction (Move 2) a dedupe store, key design
many teams react to one event a topic with one subscription per service (Move 3) copies Γ— subscribers, contract discipline
per-entity order, replay, high volume a partitioned log keyed by entity (Move 4) a partition ceiling on parallelism, hot keys, head-of-line blocking
bursts beyond capacity a buffer with headroom above the average (Move 5) waiting time
sustained overload backpressure: scale, push back, shed, prioritize (Move 5) money, errors, or lost low-value data
some messages always fail backoff + jitter, capped retries, a DLQ (Move 6) an incident queue to own; order breaks for dead-lettered keys unless you park them (parked keys stall)
consumers deploy independently additive, versioned contracts; expand/contract (contract section) slower breaking changes

Before moving on, take one design from this lesson, the checkout or the search sync, and explain it aloud as if to an interviewer, from a blank whiteboard: what stays synchronous and why, what the broker promises (at-most-once or at-least-once), how the consumer survives a duplicate and a crash, what the partition key is and what it costs, what happens during a 10Γ— burst, and where a poison message ends up. Then change one requirement (strict global order, a new team that needs 30 days of history, a provider that allows 50 calls per second) and say which part of the design moves. If you can do that, messaging questions become a checklist.

Next: fan-out at feed scale in Twitter/Instagram News Feed, a queue-driven processing pipeline in YouTube / Netflix Streaming, delivery, ordering and receipts in Chat Application (WhatsApp), and sagas and the outbox in Advanced Topics & Final Prep.