Queues & Delivery Semantics
The Warehouse System That Fell a Day Behind
Test your architecture intuition: Pitch a 7-axis solution, survive two aggressive reviewer objections, and inspect the staff-level Teacher Gold Answer.
Part 0. Start here
The problem: the order that vanished with every alarm green
Basketly sells groceries online. On Friday at about 20:00:07 the checkout of order o-7731 commits, and an OrderPlaced message for it lands on the fulfil queue. A fulfilment worker takes the message and asks the warehouse system (the WMS) for a pick list. The gift note on o-7731 has a character the WMS's label step can't handle, so the call hangs. After 40 seconds the worker gives up, logs an error and moves on. It doesn't delete the message.
The queue was set up with Amazon SQS's defaults, and two alarms that look reasonable:
- a 30-second visibility timeout (how long a taken message stays hidden) and 4-day retention, with no dead-letter queue;
- an alarm on "age of the oldest message > 5 minutes";
- and in a second version of the story (Part 5), someone adds a dead-letter queue on Saturday with an alarm on its
NumberOfMessagesSent> 0: the message moves there after 2,880 deliveries, the alarm stays green, and it still expires on Tuesday.
Nobody is paged, so nobody tells the WMS team, and the bug is never fixed.
| What happens | Arithmetic | Result |
|---|---|---|
| The message comes back every 30 s | 86,400 s a day ÷ 30 s | 2,880 deliveries a day, 11,520 in four days |
| Each delivery holds a worker slot and a WMS connection for 40 s | 40 s ÷ 30 s | 1.33 slots and 1.33 WMS connections busy with it, all the time |
| Retention deletes it | enqueue time + 4 days | Tuesday about 20:00:07: the order is gone |
Neither alarm ever fires. On a standard queue, AWS's CloudWatch metrics page says: "if a message is received three or more times and not deleted, SQS moves it to the back of the queue", and "poison-pill messages … are excluded from this metric until successfully processed". (AWS's dead-letter queue page describes the same move only for queues whose redrive policy has maxReceiveCount above 3; this page follows the metrics page.) And the dead-letter alarm watches a metric that "does not capture messages sent to a DLQ because of failed processing attempts". On Wednesday a customer asks where the groceries are. Step 1.6 of the message queue loop tells the same story in one line: "for three hours it keeps coming back and taking workers down with it".
o-7731 failed 11,520 times over four days and then vanished. Nobody was paged. Which alarm should have fired, why didn't either of ours, and what would you change so the order is set aside within minutes and done exactly once after the fix?
The big picture
Synthesizing vector architecture diagram...
What to notice: one checkout becomes one message on each of three queues, and b2b gets only the orders its filter matches. The red WMS box is the only broken thing, yet the whole fulfil path keeps paying for it. The dashed fulfil-dlq is where a failing message should go, and as shipped it doesn't exist.
What you'll be able to do after this page
- Explain why a queue lends a message rather than handing it over, and say which delivery semantics follow from where you delete (Part 1).
- Write a worker's receive loop: long polling, batch timers, the inbox check, the delete, and the in-flight limit (Part 2).
- Size a visibility timeout, extend it with a heartbeat, and explain a delete that deletes nothing (Part 3).
- Classify failures, back off with the visibility timeout, defer healthy work without burning receives, and bound retries with
maxReceiveCount(Part 4). - Run a dead-letter queue that someone notices: the right alarm metric, retention, an owner, and a safe redrive (Part 5).
- Consume a queue with Lambda without dead-lettering healthy batch-mates or losing a message, and cap its concurrency at the source (Part 6).
- Use FIFO message groups, and choose between dead-lettering, blocking and parking when one message in a group fails (Part 7).
- Read a partitioned log with a consumer group: commit after the effect, bound in-line retries, and avoid a rebalance cascade (Part 8).
- Tell the age of the oldest message from the wait of a new one, compute a drain, and scale consumers on backlog (Part 9).
- Fan out with SNS without silent drops, drain a subscription's dead-letter queue, and carry payloads over a hop's size limit (Part 10).
- Trace one order through every hop, with the ack point and the guard at each (Part 11).
- Map all of it to AWS, and name the look-alikes (Part 12).
You may have arrived from a step that relies on this: step 1.6 of the message queue loop (one bad message crashes every worker, forever), step 1.2 of the notification loop (the provider returned an error), step 2.6 of the job scheduler loop (deferral that must not burn receives), step 1.3 of the web crawler loop (one message per host group at a time), step 2.3 of the outbox ledger loop (one bad event blocks everything behind it), or the drill The Warehouse System That Fell a Day Behind. This page is the "why" behind all of them, and behind the queue and consumer lines in 22 of the 37 interview loops.
Part 1. A loan, not a hand-off
A worker took a message and died. Who has the message now? The answer decides whether the order is picked once, twice or never. This Part names the two shapes a message system can take, the one decision that sets its delivery semantics, and the example the rest of the page follows.
Queue or log
| Queue (SQS, RabbitMQ) | Log (Kafka, Amazon MSK, Kinesis Data Streams) | |
|---|---|---|
| What a consumer takes | A message, leased on receive and removed on delete | A position in a partition; reading changes nothing |
| Several consumers | Compete for messages: each message goes to one of them | Each consumer group reads every record, keeping its own position per partition |
| Fan-out to several teams | One copy per queue (SNS in front, Part 10) | Add a consumer group |
| Replay last week | Not possible: a deleted message is gone | Seek back, within retention |
| Order | FIFO queues: per message group (Part 7) | Per partition (Part 8) |
| One slow or failing message | Only that message waits (standard queue) | Every partition its consumer owns waits behind it (Part 8) |
| Scaling consumers | Add workers | At most one consumer per partition in a group |
Step 2.3 of the message queue loop makes the same choice for three teams and a replaying analytics job. Each column wins some rows: the log wins replay and cheap fan-out, and the queue wins "only this message waits". A log can't skip one bad record without an error boundary you build (Part 8); a queue can't replay.
The ack point decides the semantics
Every consumer does three things: take the message, do the work, and tell the broker it is done (a delete on a queue, an offset commit on a log). That last step is the ack point. A crash can happen anywhere, so its position decides what a crash costs.
Synthesizing vector architecture diagram...
What to notice: the same crash, 0.2 s after the receive, loses the order if the delete came first, and only delays it if the delete comes after the work. Nothing else changed.
| Where the ack is | Semantics | What a crash costs |
|---|---|---|
| Before the work | At-most-once | The message is gone, and nothing anywhere says so |
| After the work is durable | At-least-once | The work may run again |
| After the work, with an effect that ignores repeats | Effectively-once | Nothing: the repeat is recognised and skipped |
"Exactly once" exists only inside one store, such as the checkout's order row and outbox row committed in one transaction. Across a queue it is at-least-once plus an effect that ignores repeats. Idempotency & Effectively-Once Processing owns the keys and the inbox that make repeats harmless (Parts 1 and 6 there); this page owns where the ack sits and what the broker does between the receive and the delete.
Side rows 1x and 1y. A healthy order, o-7733, is enqueued at 20:00:20.65. The worker that takes it (W1 in both seeds) is killed 0.2 s later, before its WMS call returns, so the WMS creates nothing:
| # | Setting | What happened |
|---|---|---|
| 1x | Delete on receive | o-7733 was deleted at 20:00:20.65, before the work. It is never picked, and no error appears anywhere. The one or two other orders the killed worker held are lost the same way |
| 1y | Delete after the work | o-7733 becomes visible again at 20:00:50.65, 30 s after the receive; another worker picks it at 20:00:51.12 and deletes it. One pick list. Had the WMS call got through before the kill, the idempotency key would have returned the same pick list |
Why a queue at all? The drill The Warehouse System That Fell a Day Behind asks it directly. Without the queue, the checkout would call the WMS itself: a slow or broken WMS would slow or fail checkouts, and a checkout would have to retry the WMS on its own. With the queue, the checkout commits in milliseconds and the WMS works at its own pace. The price is everything on this page: the order is now picked later, at-least-once, and a failure is somebody's job to notice.
The example we follow
A real shop has many queues; that is too many to watch. So one story runs through the whole page: one order message that keeps failing. It is page 10's Basketly, on the same Friday: page 10's checkout at 20:00:07 committed with its outbox row, and this page follows the message it produced. Every trace on this page comes from running a private reference simulation of this setup (the queues, the workers, the WMS, the fulfil database, CloudWatch's metrics, and every guard as a setting), with two random seeds, not from working it out by hand. The example's numbers are small so that every delivery of one message fits on a line.
| Setting | Our example | At real scale |
|---|---|---|
| Traffic | Orders 15 a second, random arrivals. o-7731 (customer c-19) commits at about 20:00:07.05; two companions commit at 20:00:20.000: o-7733 (healthy) and o-7740 (a 300-line order) | Notification loop: about 50 messages a second; message queue loop: 50 million a day |
| Publish path | The outbox relay publishes 0.6 s after the commit; SNS delivers to each queue in 50 ms. So o-7731 reaches fulfil at 20:00:07.70, and o-7733 and o-7740 at 20:00:20.65 | Page 05's relay; the email loop's SES → SNS → SQS |
| Queues | fulfil, email, b2b: SQS standard queues, long polling (20 s) | Visibility 0 s to 12 h (default 30 s); retention 60 s to 14 days (default 4 days) |
fulfil fleet | 4 workers × 10 slots = 40. Each receive asks for up to the number of free slots (at most 10); the messages of one receive run in parallel, and each is deleted on its own success | Notification loop: Lambda; YouTube loop: long-running EC2 workers |
| One order's work | Inbox check first: if the order is already recorded, delete at once. Otherwise WMS POST /pick-lists with Idempotency-Key: <order>:pick (median 400 ms, P99 1.6 s, P99.9 about 2.5 s); record the pick and the inbox row in the fulfil DB (20 ms); DeleteMessage (10 ms). About 0.508 s in all, on average | |
| The WMS | Idempotent by key: a repeated call returns the existing pick list. For o-7731, the first call creates pick list PL-501, then the label step hangs until the caller gives up; every later call finds PL-501 and hangs again. The WMS team fixes the bug at 21:30:00, but only if someone pages them | A third-party or in-house API behind a key |
o-7740 | Rendering its packing slip takes 75 s of the worker's own time, before the WMS call | YouTube: 32 s tasks on a 120 s lease |
| Q0: as shipped | Visibility 30 s; the WMS client's timeout 40 s; on any error: log it, don't delete; no heartbeat; no dead-letter queue; alarm: age of the oldest message > 300 s | Step R2.9 of the message queue loop: "visibility set to the average" |
| Target | "Picking starts within 5 minutes of checkout for 99.9% of orders" (our SLO) |
Each Part adds one guard and replays the same message: Q1 a sized lease and a heartbeat (Part 3), Q2 counting, backoff and a dead-letter queue (Part 4), Q3 a dead-letter queue that is watched and kept (Part 5), Q4 backlog alarms that mean something (Part 9) and Q5 fan-out that can't drop (Part 10). Each setting keeps all the guards before it. Parts 6 to 8 are different: they rewind to 20:00:07 on a copy and run the same message through a Lambda consumer, a FIFO queue and a Kafka topic, because those are alternatives, not extra guards. Side rows also run on copies, so nothing they do changes the main line.
The story in eleven beats: the loan, not the hand-off (Part 1); one message's round trip (Part 2); the attempt that outlives its lease (Part 3); the message that keeps failing (Part 4); the dead-letter queue nobody hears (Part 5); the batch it drags down (Part 6); the customer it blocks (Part 7); the partitions it blocks (Part 8); the backlog behind it (Part 9); the copies it fans out to (Part 10); one order, end to end (Part 11). Each Part shows only its own events; the full table is in Part 13.
Words we use
| Word | Meaning here |
|---|---|
| Delivery | One receive of a message by a consumer |
| Receive count | How many times the message has been received (ApproximateReceiveCount in SQS) |
| Visibility timeout | How long a received message stays hidden from other consumers; the lease |
| Receipt handle | The token one receive returns; a delete or visibility change must quote it |
| Ack point | The step that tells the broker the work is done: a delete, or an offset commit |
| Poison message | A message that fails every time it is processed |
| Dead-letter queue (DLQ) | A queue where messages that failed too often are set aside |
| Redrive | Moving messages back from a DLQ to be processed again |
t1 | o-7731's first receive, 20:00:07.70 |
Events 1 to 3, as shipped (Q0)
| # | Time | Event |
|---|---|---|
| 1 | about 20:00:07.05 | o-7731's checkout commits with its outbox row: exactly once, inside one database |
| 2 | 20:00:07.70 | Relay → SNS → fulfil and email (not b2b: its filter doesn't match) |
| 3 | 20:00:07.70 (t1) | A worker receives it (W1 in seed 1; W4 in seed 2). Receive count 1; hidden until 20:00:37.70 |
Snapshot M1, t1 + 1 s (Q0)
| Lane | State |
|---|---|
o-7731 | In flight with W1; receive count 1; visible again at 20:00:37.70 |
| WMS | PL-501 created at 20:00:08.10; the call is hanging; 1 connection held for it |
| Alarms | Age of the oldest message about 1 s; nothing near 300 s |
Your worker deletes each message the moment it receives it, "so no one else gets it". What does that promise, and what does it cost on the day a deploy kills a worker?
What to remember from Part 1
- A queue lends a message; only a delete ends the loan.
- Where you put the delete (or the commit) decides at-most-once or at-least-once.
- Effectively-once is the consumer's job: an effect that ignores repeats.
Part 2. One message's round trip
At 19:59:50 Basketly is healthy: 15 orders a second arrive, and the 40 worker slots are mostly idle. Before anything fails, here is what one worker does with one message, and the three small choices in that loop that later decide whether a failure stays small.
The loop every worker runs
textWORKER LOOP (one per worker process, 10 slots) loop: n = free slots, at most 10 (n = 1 if you work one message at a time) msgs = ReceiveMessage(fulfil, MaxNumberOfMessages = n, WaitTimeSeconds = 20) for each msg, in parallel, one slot each: if the inbox already has msg.order: DeleteMessage(msg.receipt_handle); next msg start the heartbeat for msg (Q1 on, Part 3) pick = WMS: create pick list, Idempotency-Key = msg.order + ":pick" one fulfil DB transaction: insert the pick and the inbox row stop the heartbeat DeleteMessage(msg.receipt_handle) <- the ack point, after the effect
Event 4, the baseline. In the run, from 19:58 to 20:10, every receive that returned anything returned exactly one message (10,792 receives for 10,792 orders in seed 1). At 15 a second with 40 slots always waiting, a message is handed to a waiting receive the moment it arrives, so batches never form. Slots busy follow Little's law (in flight = arrival rate × time each one spends inside, page 10 Part 1): 15 a second × 0.508 s ≈ 7.6 of 40 (the run: 7.64 and 7.38 on average). Capacity is 40 slots ÷ 0.508 s ≈ 79 orders a second, so the fleet runs at about 19%.
Long polling
A receive with WaitTimeSeconds 20 waits up to 20 s for a message and returns as soon as one arrives. A receive without it (short polling) returns at once, usually empty, and samples only some of SQS's servers. At night (side row 2s, one order a minute):
| Polling | Receive calls per minute for one order |
|---|---|
| Short, every 100 ms from 4 workers | 4 × 10 × 60 = 2,400 |
| Long, 20 s | 4 workers × 3 empty 20 s returns = 12, plus 1 that returns the order: about 13 |
Every call is a billed request. Step 1.5 of the message queue loop makes the same point at its scale: 1,800 requests a second of empty polls fall to about 9.
A batch starts every timer at once
A receive can return up to 10 messages, and every message in it gets its visibility deadline at the receive instant. Side row 2b forces a batch to show what that means. A deploy restarts all four workers at 20:00:07.0. W2 is back first, at 20:00:08.00, and its first receive returns 10 messages that arrived while nobody was reading, with o-7731 fifth. This W2 variant works its batch one message at a time.
Synthesizing vector architecture diagram...
What to notice: all ten 30 s timers end together at 20:00:38. W2 is still stuck on o-7731 then, so the five messages behind it become visible and other workers take and pick them. When W2 finally reaches them at about 20:00:50, its inbox check finds them done and its deletes use old receipt handles, so they do nothing.
| Messages | Receives | Picked | What it cost |
|---|---|---|---|
b1 to b4 (before o-7731) | 1 each | 20:00:08.4 to 20:00:09.9 | Nothing |
o-7731 | 7 in the first 200 s | Never (Q0) | Its own story (Part 3) |
b6 to b10 (behind o-7731) | 2 each | 20:00:38.2 to 20:00:38.4, by other workers | 30 s late and one extra receive each. The inbox check stopped a second pick |
Receive what you'll work
The rule that follows: process a batch in parallel, or receive only as many as you'll start now. If a worker handles one message at a time, or one per FIFO group (Part 7), it should ask for 1 (MaxNumberOfMessages 1). Taking 10 and handing 9 back with visibility 0 doesn't help: each hand-back was a receive, and every receive counts toward the redrive limit of Part 4. Step 1.3 of the web crawler loop had exactly this bug before its fix; it now receives one message per slot.
Deleting. Delete each message when its own effect is durable, or collect the finished ones and send DeleteMessageBatch (up to 10) to save calls. Never delete a batch as a unit before every message in it is done.
In-flight limits
Received-but-not-deleted messages are in flight. A standard queue allows about 120,000 in flight (FIFO: 120,000 too). Over the limit, short polling gets an OverLimit error and long polling simply returns nothing. Every consumer shares that one number, and messages held by long jobs or waiting out a backoff count toward it: the crawler loop's 120,000 in-flight messages cap how many hosts it can crawl at once.
What a queue is not
A queue is not a database. SQS keeps a message at most 14 days, and you can't look one up by order number. The truth about an order lives in your database; the message is a request to act on it (step 2.3 of the web crawler loop keeps its frontier's truth in DynamoDB for this reason).
W2 took 10 messages after a restart and worked them one at a time. Nothing crashed, yet five of them were delivered twice. Why, and what should W2 have asked for?
What to remember from Part 2
- Long-poll, check the inbox, work, and delete each message when its effect is durable.
- Every message in a receive starts its visibility timer at the receive.
- If you will work one at a time, receive one at a time.
Part 3. The attempt that outlives its lease
At 20:00:37.70 two workers are hanging on the same order. Nothing crashed. o-7731's first attempt is still waiting for the WMS, its 30 s lease has run out, and SQS has handed it to someone else.
Synthesizing vector architecture diagram...
What to notice: between t1 + 30 s and t1 + 40 s two workers hold o-7731 at once, each with a WMS connection. Nobody deleted anything, and nobody will: every attempt fails the same way.
| # | Time (seed 1) | Setting | Event |
|---|---|---|---|
| 5 | 20:00:08.10 (t1 + 0.4 s) | Q0 | The WMS creates PL-501, then hangs. For the worker this is an unknown: it can't tell whether a pick list exists (Idempotency, Part 1 there) |
| 6 | 20:00:37.70 (t1 + 30 s) | Q0 | Visible again; W4 receives it (count 2). Two workers hang on one order |
| 7 | 20:00:47.70 (t1 + 40 s) | Q0 | W1's 40 s client timeout: error logged, no delete, slot freed |
| 8 | 20:01:07.70 (t1 + 60 s) | Q0 | W2 receives it (count 3). From now on it goes behind every visible message, and it is out of the age metric |
| 9 | from then on | Q0 | A delivery every 30 s, each holding a slot and a WMS connection for 40 s: 1.33 of each, for as long as it lives; 2,880 deliveries a day (the run: 120 in the first hour, 11,520 in four days, 1.333 slots on average) |
Snapshot M2, t1 + 45 s (Q0)
| Lane | State |
|---|---|
o-7731 | Receive count 2; in flight with W4 until 20:01:07.70; W1 gave up at 20:00:47.70 |
| WMS | PL-501 exists, unlabelled; 1 connection held (W4's call) |
| Alarms | Age of the oldest message 45 s |
The visibility timeout is a lease
A visibility timeout is a lease on the message: for that long, no one else gets it. Like any lease, it can expire under a worker that is still alive and still working. Leases, Fencing Tokens & Distributed Locks calls an SQS receipt handle "a lease without a fence" (Part 11 there): nothing stops the old holder's effects from landing after a new holder took over.
Sizing it
Set the visibility timeout above the slowest normal processing time, not the average. Our healthy orders take about 0.5 s (P99.9 about 2.5 s), so 30 s looks generous, until o-7740 renders a 75 s packing slip. Step R2.9 of the message queue loop lists "visibility set to the average" as a production gotcha for this reason. Q1 sets it to 60 s and adds a heartbeat, and it shortens the WMS client's timeout to 5 s, about 2 × the WMS's P99.9 (page 10's rule).
The heartbeat and its cap
textHEARTBEAT (a separate thread per message being worked; Q1 on) every 20 s while the work is running: ChangeMessageVisibility(msg.receipt_handle, 60 s) -> hidden until now + 60 s stop when the work finishes, fails or is abandoned if the call is refused (the 12-hour cap), stop working: the message will come back
The new timeout counts from the call, not from the receive, and it isn't added to what was left. It is capped: "the visibility timeout has a maximum limit of 12 hours from when the message is first received. Extending the timeout doesn't reset this 12-hour limit" (SQS). A job that can run longer than 12 hours can't keep its message hidden: claim the job in your own table, delete the message, and keep your own lease (step 2.5 of the job scheduler loop). Step 2.2 of the YouTube loop uses 120 s extended every 30 s.
The delete that deletes nothing
Here is o-7740 (enqueued 20:00:20.65, a 75 s render) in three settings. The inbox check runs first on every delivery.
| # | Setting | Renders | Deliveries | What happened (seed 1; seed 2 within 0.05 s) |
|---|---|---|---|---|
| 10 | Q0: 30 s, no heartbeat | 3 | 4 | Received at 20:00:20.65, again at 20:00:50.65 and 20:01:20.65: three renders running. The first finishes and records the pick at 20:01:36.26, and its delete, with the first receipt handle, does nothing. The 4th delivery at 20:01:50.65 finds the inbox row and deletes at 20:01:50.66. The other two renders finish at 20:02:05.9 and 20:02:36.1, find the inbox row and record nothing. One pick list |
| 11a | Q1 without the heartbeat: 60 s | 2 | 3 | The second receive at 20:01:20.65 starts a second render; the first delete is a no-op again; the 3rd delivery deletes it at 20:02:20.66 |
| 11 | Q1: 60 s, heartbeat every 20 s | 1 | 1 | Extended at 20:00:40.65, 20:01:00.65 and 20:01:20.65; picked at 20:01:36.26 and deleted with a current handle at 20:01:36.27 |
AWS on a delete with an old receipt handle: "the request will succeed, but the message might not be deleted". Our simulation takes the unsafe reading: it stays. So never rely on the delete to prevent a second effect. Two guards did that here: the inbox check first, which is what finally deleted o-7740, and the WMS's key, which would have returned the same pick list if a second render had reached it.
Slow failures hold resources
o-7731 in Q0 fails slowly: 40 s per attempt, repeated every 30 s, so it holds 40 ÷ 30 ≈ 1.33 worker slots and 1.33 WMS connections for four days. Those slots and connections are shared with every healthy order. With Q1, each attempt fails in 5 s and then waits out the rest of its 60 s lease: 5 ÷ 60 ≈ 0.08 slots and connections, and 86,400 ÷ 60 = 1,440 deliveries a day (the run: 1,441 in 24 hours, 0.083 slots). Cheaper, but still forever.
Letting go early
A worker that knows it can't do the work (it is shutting down, or the message is for a host it may not touch yet) should release the message at once with ChangeMessageVisibility 0, or a short delay, instead of letting the lease run out. Step 2.5 of the web crawler loop returns unstarted messages with visibility 1 s when a Spot instance gets its interruption notice.
o-7740 was rendered three times though no worker crashed, and two deletes reported success without deleting it. What finally removed it, and what one change renders it once?
What to remember from Part 3
- Size the lease for the slowest normal job and extend it while you work.
- A lease can expire under a live worker: guard the effect (inbox first, a key), not the delete.
- A message that fails slowly holds a worker and a downstream connection every time it comes back.
Part 4. The message that keeps failing
With Q1, o-7731 fails in 5 s and costs 0.08 slots. It still fails 1,440 times a day, and nothing will ever stop it. Retrying is right for failures that can succeed later. For this one it only burns WMS calls until retention deletes the order.
Which errors to retry
Retries, Timeouts, Backpressure & Load Shedding owns the classification (Part 3 there). For a queue consumer it becomes:
| Class | Examples | What the consumer does |
|---|---|---|
| Transient | A timeout, a connection reset, 500, 502, 503, 504 | Retry later, with backoff, through the queue |
| Throttle | 429, often with Retry-After | Retry after the server's hint |
| Permanent | 400, 404, 422, a schema or validation error | Don't retry: send it to the DLQ now, then delete |
| Unknown | A timeout after the call may have acted (PL-501) | Retry with the same idempotency key |
Count, and back off with the visibility
Every receive raises the message's ApproximateReceiveCount, and the consumer sees it. Q2 uses it twice: to set a backoff and to decide when to give up.
textON A FAILED ATTEMPT (Q2) n = msg.ApproximateReceiveCount if the error is permanent: SendMessage(fulfil-dlq, msg.body) first send it... DeleteMessage(msg.receipt_handle) ...then delete (a crash between = a harmless copy) else if the error is a throttle with Retry-After: ChangeMessageVisibility(msg.receipt_handle, Retry-After) else: ChangeMessageVisibility(msg.receipt_handle, random(0, min(900 s, 10 s x 2^(n-1)))) if n = maxReceiveCount: emit LastAttemptFailed (Q3 on, Part 5)
The backoff is page 10's full jitter (Part 4 there), applied by hiding the message: after the 1st failure it waits up to 10 s, after the 2nd up to 20 s, then 40, 80 and 160 s, capped at 900 s. Step 1.2 of the notification loop uses the same method with 10 receives, and about 51 minutes to the DLQ.
Bound it: maxReceiveCountand the dead-letter queue
A redrive policy on fulfil names a dead-letter queue (fulfil-dlq) and maxReceiveCount 5: "the number of times a consumer can receive a message from a source queue before it is moved" (SQS). So o-7731 is delivered five times, and it is moved by the next receive after the count reaches 5. Our workers always have a receive open, so that is the moment its 5th visibility timeout ends. If nothing is receiving, nothing is moved (Part 5).
Synthesizing vector architecture diagram...
What to notice: the only way into the DLQ is through a receive. The move is triggered by a consumer asking for work, not by a timer, so a queue nobody polls never moves anything.
A DLQ must be the same type as its source (a FIFO queue needs a FIFO DLQ), in the same account and Region.
Event 12 (Q2, seed 1). Five deliveries, each failing in 5 s, with jittered waits after them:
| Delivery | Received | Failed | Backoff drawn (cap) | Visible again |
|---|---|---|---|---|
| 1 | 20:00:07.70 | 20:00:12.70 | 2.8 s (10 s) | 20:00:15.46 |
| 2 | 20:00:15.46 | 20:00:20.46 | 2.0 s (20 s) | 20:00:22.45 |
| 3 | 20:00:22.45 | 20:00:27.45 | 34.7 s (40 s) | 20:01:02.18 |
| 4 | 20:01:02.18 | 20:01:07.18 | 40.2 s (80 s) | 20:01:47.41 |
| 5 | 20:01:47.41 | 20:01:52.41 | 24.2 s (160 s) | 20:02:16.60: moved to fulfil-dlq |
Snapshot M4, at the move (seed 1, 20:02:16.60)
| Lane | State |
|---|---|
o-7731 | In fulfil-dlq; it was received 5 times |
| WMS | PL-501 exists, unlabelled; no connections held |
| Alarms | See Part 5: in Q2 there is still no alarm on the DLQ |
Seed 2 moves it at 20:03:29.33 (t1 + 201.6 s). How long does it take in general? Each of the five attempts takes 5 s, and the jittered waits average 5 + 10 + 20 + 40 + 80 = 155 s, so about 180 s on average, anywhere from 25 s to 335 s. Over 5,000 seeds the run gives a mean of 179.7 s, a range of 38.7 to 318.6 s, and 91.4 to 268.7 s between the 5th and 95th percentiles. The 5th failure comes at a mean of 100.0 s, about 80 s before the move.
| # | Setting | WMS label attempts for o-7731 |
|---|---|---|
| 13 | Q0, four days | 11,520 |
| 13 | Q2 | 5 |
Deferring healthy work
Now a different failure: nothing is wrong with any order, but the WMS throttles everyone. Replay T (a copy of Q2): from 20:10:00 to 20:20:00 every WMS call returns 429 with Retry-After: 60. Four ways to wait it out (seeds 1 and 2):
| Deferral | Healthy orders in the DLQ | Longest wait, arrival to pick | What the age alarm shows |
|---|---|---|---|
Visibility set to Retry-After, 60 s | 5,374 and 5,421 | 243 s; 245 s for the ones that made it | At most 120 s |
| Q2's jitter instead | 7,904 and 8,017 | 138 s; 143 s | At most 29 s |
Delete, then re-send with DelaySeconds 60 | 0 | 655 s; 658 s | At most 120 s: each re-send is a new message whose age starts at zero |
| Pause the consumer from 20:10 to 20:20 | 0 | 601 s (both seeds) | 601 s: the alarm sees it |
Why 5,400? Each deferral is a receive, so an order that arrives before 20:20:00 − 4 × 60 s = 20:16:00 has all five of its deliveries inside the throttle and is moved: 360 s × 15 a second = 5,400 (the run: 5,374 and 5,421, all enqueued 20:10:00 to 20:16:00). Jitter makes it worse because its first waits are shorter than 60 s.
| On equal terms | Visibility (backoff) | Delete + re-send with a delay | Pause the consumer | A scheduler (EventBridge Scheduler) |
|---|---|---|---|---|
| Counts a receive? | Yes | No: a new message | No | No |
| Longest delay | 12 h from the message's first receive | 15 minutes | As long as you pause | Any |
| FIFO queues | Yes | No per-message delay on FIFO | Yes | Yes |
| Order | Kept on FIFO | Lost: it rejoins at the back | Kept | Lost |
| Age metric | Sees it only until the 3rd receive, then leaves it out (replay T: at most 120 s while orders waited 243 s) | Resets it (at most 120 s while orders waited 655 s) | Sees it | Sees nothing |
| Can the DLQ still catch a real poison message? | Yes, too eagerly | Only if you count attempts yourself (in a message attribute) | Nothing moves while paused | Only in your own code |
Use the visibility timeout for real retries of a message that failed. Defer healthy work by re-sending it with a delay (carrying your own attempt count), or by pausing the consumer, as step R1.9 of the notification loop does during a long provider outage. Step 2.6 of the job scheduler loop says it plainly: "every receive counts toward the queue's maxReceiveCount".
Delay queues and timers
DelaySeconds hides a new message for up to 15 minutes, set per queue or, on standard queues only, per message (FIFO queues don't support per-message timers). Step 1.4 of the ride-sharing loop uses a 16 s delayed message as a durable timer, at-least-once, guarded by a conditional update. For longer waits, SQS's documentation points to EventBridge Scheduler, which is a scheduler, not a queue (Part 12).
The WMS throttled everyone for 10 minutes, and nothing was wrong with any order. Why are 5,400 of them in the DLQ?
What to remember from Part 4
- Count every delivery; after N, move the message aside instead of retrying forever.
- Back off real failures with the visibility timeout; defer healthy work by re-sending with a delay or pausing.
- A permanent error goes straight to the DLQ.
Part 5. The dead-letter queue nobody hears
Back to the hook. In Q0, o-7731 is delivered 11,520 times from Friday to Tuesday. The age alarm stays green, because from its third receive the message is left out of the metric. On Tuesday at 20:00:07.70 retention deletes it. Nobody was paged; on Wednesday a customer asks.
Snapshot M3, Saturday 20:00 (Q0)
Synthesizing vector architecture diagram...
What to notice: in "Queue" the message has been received 2,880 times in its first day; in "Workers and WMS" it holds a slot and a connection all the time; in "Alarms" everything is green, because the age metric doesn't count it.
Side row 5a (a copy of Q0, still no fix). On Saturday at 20:00 someone adds fulfil-dlq (default 4-day retention, maxReceiveCount 5) and an alarm on its NumberOfMessagesSent > 0. o-7731's count is already 2,880, so it moves at its next receive, Saturday 20:00:07.70. The alarm never fires: the move isn't a send. The DLQ's own age metric starts at the move and reads 24 h on Sunday, 48 h on Monday and 72 h at 19:59 on Tuesday, while retention still counts from Friday: the message expires on Tuesday at 20:00:07.70, 3.00 days after it arrived. The age metric makes it look a day younger than its expiry clock says.
Five ways a DLQ fails silently
| Failure | What goes wrong | The fix |
|---|---|---|
| No alarm | Messages pile up where nobody looks | An alarm on every DLQ |
| An alarm that can't fire | NumberOfMessagesSent ignores redrive-policy moves; an idle DLQ's first datapoint can be up to 15 minutes late ("a delay of up to 15 minutes occurs … when a queue is activated from an inactive state") | Alarm on ApproximateNumberOfMessagesVisible ≥ 1, and emit your own LastAttemptFailed metric when a message with receive count = maxReceiveCount fails |
| Retention that expires it | On a standard queue a moved message keeps its original enqueue time; FIFO resets it | DLQ retention longer than the source's: 14 days, as AWS advises |
| No owner | The alarm fires into a channel nobody reads | An on-call owner, a runbook, and a page to whoever can fix the cause |
| Nobody polling | The move happens at a receive; a queue with its consumers scaled to zero (or its Lambda mapping disabled) moves nothing | Know that "no DLQ messages" can mean "nothing is receiving"; alarm on the age of the source too |
The age metric's blind spot
The hook's quote again: "if a message is received three or more times and not deleted, SQS moves it to the back of the queue … This reordering occurs even when a redrive policy is in place", and such messages are "excluded from this metric until successfully processed" (the CloudWatch metrics page). AWS's DLQ page describes the reordering only for standard queues whose maxReceiveCount is above 3, so treat the exact condition as a detail that may change; the lesson doesn't. The age alarm catches backlogs, not one poison message. The receive count, the DLQ and your own last-attempt metric catch that. FIFO queues don't reorder, so there the age alarm does see a stuck message (Part 7).
An alarm that fires (Q3)
Q3 keeps Q2's policy and adds: fulfil-dlq retention 14 days; an alarm on its ApproximateNumberOfMessagesVisible ≥ 1; LastAttemptFailed; and an owner, the fulfil on-call, who pages the WMS team. CloudWatch is modelled as 1-minute datapoints (the maximum over the minute, published at the minute's end), with each alarm firing on its first breaching datapoint.
| # | Time (seed 1) | Event |
|---|---|---|
| 15 | 20:01:52.41 | The 5th failure raises LastAttemptFailed; its alarm fires at 20:02:00, before the move |
| 15 | 20:02:16.60 | o-7731 moves to fulfil-dlq (t1 + 128.9 s). The DLQ-depth alarm fires at 20:03:00, or at 20:18:00 if CloudWatch takes its full 15 minutes to report a DLQ that was idle all day |
| 16 | 21:30:00 | Paged, the WMS team deploys the label fix |
| 17 | 21:40:00 | Redrive fulfil-dlq → fulfil at 10 messages a second: o-7731 is re-enqueued as a new message (new ID, new enqueue time, receive count 0) and received at once. The WMS returns PL-501 by its key, the label step now works, and the pick is recorded at 21:40:00.52. Deleted at 21:40:00.53. Picking starts 1 h 40 min late, with one pick list |
In seed 2 the 5th failure is at 20:01:55.91 (alarm 20:02:00), the move at 20:03:29.33 (DLQ alarm 20:04:00 or 20:19:00), and the pick at 21:40:00.68.
Snapshot M5, 21:40:01 (Q3), next to M3
Synthesizing vector architecture diagram...
What to notice: the same three lanes as M3. In "Queue" the message is gone and the DLQ is empty; in "Workers and WMS" one pick list exists, found by its key; in "Alarms" both alarms fired within a minute of each other, which is what got the bug fixed.
Bringing it back
A redrive (StartMessageMoveTask) moves messages from a DLQ back to their source, or to another queue. Four facts shape how you use it:
- A redriven message is a new message: new ID, new enqueue time, receive count 0. If it fails again, it goes through all five deliveries again.
- It interleaves with new traffic; order is not restored.
- A custom velocity of up to 500 messages a second; a task runs at most 36 hours. (A message moved into a FIFO DLQ has its deduplication ID replaced by its original message ID.)
- It works only for DLQs whose sources are SQS queues, not for an SNS subscription's DLQ (Part 10).
textREDRIVE RUNBOOK (fulfil-dlq) 1. look: read a sample; group by error; confirm it is one cause 2. fix: deploy the fix, or get the owning team to (the WMS label step) 3. test: redrive ONE message; watch it succeed end to end 4. move: StartMessageMoveTask at a capped velocity (10/s here; far below the fleet's 79/s) 5. watch: the DLQ draining, the source's age and errors; stop the task if failures return 6. note: what failed, how many, when fixed, when redriven
Step R3.9 of the message queue loop uses the same six steps at 50 a second.
Side row 5c (a copy of Q3). Someone redrives at 21:00, before the 21:30 fix. o-7731 becomes a new message with a count of 0, hangs five more times, and is back in the DLQ at 21:04:05 (seed 2: 21:03:25). Nothing was gained, and the WMS took five more hung calls. Look before you move.
Repair instead of replay. Sometimes the message is only a pointer, or the order has changed since it was sent. Then re-read the current state from the source and act on that, instead of redriving the old message: Change Streams & the Transactional Outbox owns this (Part 6 there).
Your DLQ alarm is NumberOfMessagesSent > 0, and the DLQ holds 30 messages. Why did the alarm never fire, and when will those messages disappear?
What to remember from Part 5
- A DLQ needs an owner, an alarm on its visible depth, and longer retention than its source.
- On a standard queue the age metric can't see a poison message; count receives and alarm on the last failure.
- Redrive only after the fix, at a capped rate: every redriven message is new.
Part 6. The batch it drags down
This Part rewinds to 20:00:07 on a copy of Q2 (the same queue and fulfil-dlq, maxReceiveCount 5) and replaces the worker fleet with a Lambda function. Ten orders end up in the dead-letter queue, and only one of them was ever broken.
How a Lambda reads a queue
A Lambda event source mapping polls the queue for you, gathers messages into a batch and invokes your function with the whole batch. You choose:
- the batch size: up to 10,000 for a standard queue, 10 for FIFO; a size above 10 needs a batching window of at least 1 s;
- the batching window: how long to gather messages before invoking, up to 5 minutes, on standard queues only.
When messages are available, Lambda starts with 5 concurrent batches and adds up to 300 more concurrent invokes a minute, up to 1,250 per mapping. (Provisioned mode instead keeps a set of dedicated pollers, each handling up to 1 MB a second and 10 concurrent invokes, and scales faster.) The mapping deletes a batch's messages when the function succeeds. When the invoke fails as a whole, it deletes none of them.
Timeouts that fit. AWS recommends a visibility timeout of at least "six times your function timeout, plus the value of MaximumBatchingWindowInSeconds", so that Lambda can retry a throttled batch before the messages come back. And the function's own timeout must cover a whole batch, including one hung call: 9 healthy orders × about 0.5 s + one 5 s WMS timeout ≈ 9.6 s. With a 10 s function timeout, that batch timed out in 25.5% of 1,000 seeds in the run, so replay L uses 15 s.
| Replay L setting | Value |
|---|---|
| Batch size, window | 10, 1 s |
| Function timeout | 15 s |
| Visibility timeout | 100 s (6 × 15 s + 1 s = 91 s, rounded up) |
| WMS timeout | 5 s |
| The batch | o-7731 arrives 5th in a batch of 10. That is our choice: at 15 a second, spread over 5 concurrent pollers, real batches here hold 1 to 3 messages |
| The function | Handles its records one after another, in order, inbox check first |
Synthesizing vector architecture diagram...
What to notice: the same batch of 10 with the same broken order, and three handlers. "Throws" sends all ten to the DLQ; "Catches, returns an empty list" loses the broken order silently; only "Returns o-7731's ID" sets aside exactly the one message that failed.
A failed batch comes back whole
| # | Handler (seeds 1 and 2) | What happened |
|---|---|---|
| 19 | Throws at o-7731 | Invoked at 20:00:08.00. It handles m1 to m4, then hangs 5 s on o-7731 and throws. The batch's ten timers expire together, so the ten come back together and the batch re-forms: invokes at 20:00:08, 20:01:48, 20:03:28, 20:05:08 and 20:06:48. m1 to m4 are handled 5 times each (16 repeat runs, each stopped by the inbox check; without it, 16 duplicate picks for the WMS key to catch). m6 to m10 are never handled. At 20:08:28 all 10 are moved to fulfil-dlq |
| 20 | Catches, returns {"batchItemFailures": []} | One invoke; every mate is picked, and o-7731 is caught and forgotten. An empty list means complete success, so Lambda deletes all 10 at 20:00:18.1 (seed 2: 20:00:17.4). o-7731 is lost: not in any queue, not in the DLQ, and the deleted-messages metric looks healthy |
| 21 | Returns o-7731's message ID | The 9 mates are picked once and deleted. Only o-7731 comes back, every 100 s, alone; it is moved to the DLQ at 20:08:28 after 5 deliveries. 0 repeat runs |
Where the broken message sits in the batch matters when the handler throws: every mate before it is handled 5 times. Averaged over all ten positions, that is 18 repeat runs per batch; if SQS returns the re-formed batch in a shuffled order, 13.8 (the run: 200 seeds per position). AWS re:Post has the symptom in its title: "Lambda retrying valid SQS messages and placing them in my DLQ".
Report failures by ID
textBATCH HANDLER WITH ReportBatchItemFailures ON failures = [] for each record in the batch, in order: try: process(record) inbox first, WMS with the key, DB row catch any error: failures.append(record.messageId) on a FIFO queue: also append every record after it, then stop return batchItemFailures = one {"itemIdentifier": id} per id in failures
| Lambda treats the batch as... | When the function returns... |
|---|---|
| A complete success (deletes every message) | An empty batchItemFailures list, a null list, an empty or null response |
| A complete failure (deletes nothing) | It throws; invalid JSON; an empty or null itemIdentifier |
| A partial success | A list of the failed message IDs: only those come back |
On a FIFO queue, AWS says the function "should stop processing messages after the first failure and return all failed and unprocessed messages", or the group's order breaks. With partial batch responses on, Lambda also doesn't slow its polling down when invokes fail.
Concurrency at the source, not at the function
To stop a burst from flooding the WMS, cap the event source mapping's maximum concurrency (2 to 1,000), not the function's reserved concurrency. When reserved concurrency throttles an invoke, the batch goes back to the queue with every message's receive count raised, and Lambda keeps retrying "until the message's timestamp exceeds your queue's visibility timeout". Healthy messages can reach the DLQ that way: in the AWS Compute Blog's demonstration, throttling sent 9 messages to the DLQ, and maximum concurrency sent none. Keep maximum concurrency at or below the reserved concurrency, and maxReceiveCount at least 5, as AWS advises. (Lambda's account concurrency, 1,000 by default, is shared by every function in the Region: an uncapped mapping can take all of it.)
Kinesis and DynamoDB Streams sources work differently: records are read by position, and a failing batch blocks its shard. Change Streams & the Transactional Outbox owns bisecting a batch, the maximum record age and the on-failure destination (Part 9 there).
You turned on ReportBatchItemFailures, and your handler returns {"batchItemFailures": []} when it catches an error. What happens to the failed message?
What to remember from Part 6
- One bad message fails its whole batch, and healthy mates share its fate, up to the DLQ, unless you report it by ID.
- An empty failure list means "all succeeded": never catch and return nothing.
- Cap concurrency at the event source, not by throttling the function.
Part 7. The customer it blocks
This Part rewinds to 20:00:07 on a copy of Q2's policies, applied to a FIFO queue fulfil.fifo with MessageGroupId = the customer. Customer c-19 sends three messages: o-7731 (20:00:07.70), "change the delivery slot for o-7731" (20:00:40.65), and a new order o-7802 (20:01:05.65). About 5,000 other customers keep ordering at 15 a second. Without a dead-letter queue, o-7802 is picked on Tuesday.
Message groups
A FIFO queue keeps order within a message group, and only there. While a group has a message in flight, "subsequent messages in that group are not made available", so a group is worked one receive's worth at a time; different groups run in parallel. A receive "attempts to return as many messages as possible with the same message group ID", which is why our workers ask for one message per receive (Part 2): taking a group's next nine messages and handing them back would burn their receives.
Synthesizing vector architecture diagram...
What to notice: order is kept only inside each group. c-19's two later messages (yellow) wait behind o-7731 (red); the other customers' groups don't wait for anything.
Head-of-line blocking, and why the age alarm works here
| # | Time (seed 1) | Setting | Event |
|---|---|---|---|
| 22 | 20:00:07.70 on | F0 | o-7731 fails and backs off as in Q2; c-19's group has it in flight the whole time, so the slot change and o-7802 wait. Every other group flows (28,881 other orders picked while the run models traffic, 19:58 to 20:30) |
| 23 | Friday → Tuesday | F1: no DLQ | o-7731 is delivered 786 times (seed 2: 730) until retention deletes it on Tuesday at 20:00:07.70. Then the slot change runs against an order fulfil never picked: "order unknown", which we treat as retryable, so it fails until its own retention ends on Tuesday at 20:00:40.65. o-7802 is picked on Tuesday at 20:00:41.36, four days late. Each later message keeps its own retention clock, so it too could have expired |
| 23 | 20:02:07.70 → 20:03:00 | F1 | The age alarm does fire: a FIFO queue doesn't move a message to the back, so o-7731 stays the oldest. Its age passes 120 s at 20:02:07.70, and Q4's 120 s alarm fires at the 20:03:00 datapoint (Q2's 300 s alarm: 20:06:00) |
A DLQ breaks the order; parking keeps it
| # | Setting (seed 1; seed 2 in brackets) | What happened | Order after the repair |
|---|---|---|---|
| 24 | F2: a FIFO DLQ, maxReceiveCount 5, with Q3's page and the 21:30 fix | o-7731 moves at 20:02:16.60 (20:03:29.33). The slot change runs next, against an order fulfil has never seen: it fails 5 times and is moved at 20:05:48.28 (20:07:12.31). o-7802 is picked at 20:05:48.68 (20:07:12.79). After the fix, a redrive at 21:40 replays the DLQ in its order: o-7731 picked at 21:40:00.84, the slot change applied at 21:40:00.87 | Broken: o-7802 was picked 1 h 34 min before the order ahead of it, and the slot change needed 5 useless attempts |
| 25 | F3: block and park, with Q3's page and the 21:30 fix | At the 5th failure, 20:01:52.41 (20:01:55.91), the consumer records o-7731 in its own dead_letters table, marks c-19 blocked and deletes the message. The slot change and o-7802 arrive at once; each is parked in order in the consumer's table and deleted. At 21:40 a repair applies all three in order: o-7731, its slot change, o-7802 | Kept: the group's order survives, and every other customer kept flowing |
AWS says it directly: "Don't use a dead-letter queue with a FIFO queue if you don't want to break the exact order". Step R2.8 of the message queue loop weighs "move to DLQ and continue" against "block the group and page". F3 is a third answer: page 05's error boundary (Change Streams & the Transactional Outbox, Part 6 there), applied to a queue.
| On equal terms | DLQ and continue (F2) | Block and page (F1 + an age alarm) | Block and park (F3) |
|---|---|---|---|
| The group's order | Broken | Kept | Kept |
| The group's later messages | Processed at once, possibly against missing state | Wait until someone fixes it | Recorded, not processed |
| Other groups | Flow | Flow | Flow |
| Operator effort | Redrive, and reconcile anything processed out of order | Fix before retention runs out | Fix, then run the repair from the park table |
| What you must build | Nothing | An alarm that catches it | A dead_letters table, a blocked-keys set, a park table and a repair job |
Snapshot M6, 20:02:30 (seed 1)
| F1 (no DLQ) | F3 (block and park) | |
|---|---|---|
o-7731 | In flight with its 6th delivery, backing off | In dead_letters; message deleted |
Slot change, o-7802 | Waiting in the group, never received | Parked in order; messages deleted |
Group c-19 | Blocked | Blocked in the consumer's table |
| Other orders picked by then | 3,985 | 3,985 |
Choosing the key
Side row 7k (a copy of F1, first hour). Group by the smallest thing whose order matters, no wider:
MessageGroupId | What waits behind o-7731 |
|---|---|
| The customer | c-19's 2 later messages, for 4 days |
| The order | Only the slot change for o-7731; o-7802 is its own group and is picked at 20:01:05.88 |
| One group for the whole queue | Everything. And one group can't keep up even without a poison message: one receive's worth in flight at a time is about 1 ÷ 0.508 s ≈ 2 orders a second against 15 arriving. In the first hour only 1,679 orders were picked, all of them from the backlog that was already ahead of o-7731 |
A hot group, like a hot partition key, is a hot key (Sharding, Hot Keys & Rebalancing).
Deduplication is not exactly-once. A FIFO queue drops a repeated send with the same deduplication ID for 5 minutes. It still redelivers a message whose visibility timeout ran out, as Part 3 showed. Idempotency owns this (Part 6 there), and step 1.3 of the notification loop corrects "use an SQS FIFO queue for exactly-once".
Throughput
| FIFO limit | Value |
|---|---|
| Per API action (send, receive, delete), each on its own | 300 transactions a second; 3,000 messages a second with batches of 10 |
| High-throughput mode, per action | 70,000 transactions a second in N. Virginia, Oregon and Ireland (700,000 messages a second batched); 19,000 in Ohio and Frankfurt; lower elsewhere, down to 2,400 |
| Heartbeats | Each ChangeMessageVisibility is one more API call; AWS advises minimizing "frequent visibility changes" to save the queue's TPS |
| Lambda on FIFO | Concurrent invokes ≤ the number of active message groups (and ≤ maximum concurrency); batch size ≤ 10; no batching window |
The transactions a second are shared by every producer and consumer of the queue.
Fair queues: tenants without order
Steps 3.1 of the notification loop and 3.2 of the job scheduler loop face the opposite problem: one tenant's flood delays everyone, and nobody needs order. SQS fair queues solve it on a standard queue: set MessageGroupId to the tenant, and when one tenant has "a disproportionately large number of in-flight messages", SQS "prioritizes message delivery for other tenants". A fair queue keeps no order, and it "does not limit the consumption rate per tenant".
Would you put a DLQ on the payments FIFO queue? What happens to the capture that follows a dead-lettered authorize?
What to remember from Part 7
- FIFO orders per group; one failing message stops its whole group.
- A DLQ unblocks the group but breaks its order; parking the key keeps both.
- Group by the smallest thing whose order matters.
Part 8. The partitions it blocks
This Part rewinds to 20:00:07 on a copy of Q1's client (a 5 s WMS timeout, no backoff) reading a Kafka topic instead of a queue. One record stops a third of all customers at once, and ten minutes later nothing reads the topic at all.
| Replay K setting | Value |
|---|---|
Topic orders | 6 partitions, key = the customer; c-19 is on partition 2 |
Consumer group fulfil | 3 consumers, C1 < C2 < C3. The default assignor list is [RangeAssignor, CooperativeStickyAssignor], so the first, Range (eager), is used: C1 gets partitions 0 and 1, C2 gets 2 and 3, C3 gets 4 and 5 |
| Each consumer | One poll() loop. It works the records of one poll with different keys in parallel and each key's records in offset order, and polls again when they are all done |
| Offsets | Auto-commit (enable.auto.commit true) every 5 s (auto.commit.interval.ms 5,000); max.poll.records 500 |
| Failure detection | max.poll.interval.ms 300,000 (5 minutes); session.timeout.ms 45,000 |
| A rebalance | Pauses the partitions it moves for 3 s (our figure): with Range, all of them |
| Retry policy (K1) | Retry the failing record in line, forever |
Consumer groups and offsets
In a group, each partition is read by one consumer, so a group can't use more consumers than partitions. A consumer's position in each partition is an offset, and committing it is the log's ack point: "the committed offset should always be the offset of the next message that your application will read" (the KafkaConsumer documentation). Commit after the effect is durable, and it is at-least-once; before, at-most-once. Auto-commit commits, at each poll(), the position of the records the previous poll() returned. That is at-least-once only if every record is finished before the next poll(), which our loop does.
Synthesizing vector architecture diagram...
What to notice: o-7731 is on partition 2, but partition 3 stops too, because C2's one poll() loop serves both. Two of six partitions, a third of all customers, wait behind one record.
One record stops a consumer's partitions, then the group
| # | Time (seed 1) | Event |
|---|---|---|
| 27 | 20:00:08.48 | C2 starts o-7731 and retries it forever, 5 s per attempt. It never returns to poll(), so partitions 2 and 3 stop. Lag grows at 15 ÷ 3 = 5 records a second (the run: 5.1) |
| 28 | 20:05:08.48 | No poll() for 300 s: C2 "is considered failed and the group will rebalance". The eager rebalance pauses all 6 partitions for 3 s |
| 28 | 20:05:11.48 | With two members, Range gives C1 partitions 0, 1 and 2, and C3 partitions 3, 4 and 5. C1 reads partition 2 from C2's last committed offset, re-reading the records C2 had finished since that commit, and reaches o-7731 at once. Now half the customers stop: lag grows at 7.5 a second (the run: 7.7) |
| 28 | 20:10:11.48 | C1 is considered failed; another 3 s pause; C3 owns all 6 partitions from 20:10:14.48 |
| 28 | 20:10:25.22 | C3 works through partitions 0 and 1's backlog first, then reaches o-7731: every partition stops; lag grows at 15 a second |
| 28 | 20:15:25.22 | C3 is considered failed. The group is empty: nothing reads the topic. At 20:24 the lag is 14,031 records and still growing at 15 a second |
Synthesizing vector architecture diagram...
What to notice: each consumer that inherits partition 2 meets o-7731 and stops within seconds; the red bars follow each other until the group is empty. Stopped partitions grow the lag by 5, then 7.5, then 15 records a second.
The lag metric told the truth throughout: Amazon MSK emits consumer-lag metrics (MaxOffsetLag, SumOffsetLag, EstimatedMaxTimeLag and others) only while a group is STABLE or EMPTY, and here that was every second except the three 3 s rebalances (the third, at 20:15:25, leaves the group empty). A group that never settles, such as consumers crash-looping, has no lag datapoints at all, so alarm on missing data too. The run also shows the cost of each rebalance: 432 healthy records were read twice, and the inbox kept every pick single.
Side row 8a: auto-commit with a thread pool. A consumer that hands each polled record to a pool of 8 threads and polls again at once commits, at its next poll(), the position of records it has only handed off. Kafka's documentation warns that the committed offset can then "get ahead of the consumed position". On a copy, C2 does exactly this; o-7731 and its slot change occupy two pool threads, retrying. C2 is killed at 20:02:00, the group notices after the 45 s session timeout (20:02:45) and gives partition 2 to C1 from the committed offset, which is past o-7731 and the slot change. Both are skipped for good, silently. And the pool had already started o-7802 at 20:01:05.70, ahead of the order in front of it.
Rebalances
Two settings stop a restart from freezing the group (step 2.6 of the message queue loop):
- The cooperative assignor (
CooperativeStickyAssignor, listed alone): a rebalance moves only the partitions that change owner, and the others keep flowing. The default list uses Range first, which revokes everything. - Static membership (
group.instance.id): a consumer that restarts withinsession.timeout.msgets its partitions back without a rebalance at all. Sizesession.timeout.msabove a normal restart (45 s by default), so a routine deploy doesn't rebalance; the price is that a consumer that really died holds its partitions that long.
Neither helps a consumer that never calls poll(). That needs the next fix.
The fix: bound the retries, then page 05's boundary
textCONSUMER LOOP WITH AN ERROR BOUNDARY (K2) records = poll() for each key in records, in offset order (keys in parallel): if key is blocked: park the record in the fulfil DB; next for attempt in 1..4: waits 0.1 s, 1 s, 5 s between try: process(record); break catch: if attempt = 4: record it in dead_letters (fulfil DB); block key commit the offsets after the whole poll is processed (cooperative-sticky assignor, group.instance.id set)
This is Change Streams & the Transactional Outbox's error boundary (Part 6 there): retry a few times, dead-letter in the consumer's own database, block the key, park its later records, and keep committing.
| # | Time (seed 1) | K2 |
|---|---|---|
| 29 | 20:00:08.48 → 20:00:34.60 | Four attempts of 5 s with waits of 0.1, 1 and 5 s between: 4 × 5 + 6.1 = 26.1 s. Then o-7731 is dead-lettered and c-19 blocked |
| 29 | during those 26 s | Partitions 2 and 3 wait once: the longest wait of a record on them is 26.0 s and 25.9 s (seed 2: 25.7 and 25.8), against 2.4 s on the other partitions |
| 29 | 20:00:40.84, 20:01:05.83 | The slot change and o-7802 are parked in order; no rebalance; all 3 consumers stay |
Retry topics (send the failing record to a retry topic and move on) also unblock the partition, but they give up the key's order, like a DLQ on a FIFO queue.
The same poison, four transports
| On equal terms | Standard queue (Q2) | FIFO queue (F) | Kafka partition (K) | Lambda batch (L) |
|---|---|---|---|---|
| Who else waits | Nobody | Its group | Every partition its consumer owns; then more after each rebalance | Its batch-mates |
| By default, for how long | Nobody waits, but it retries until retention (4 days) | Until it expires (4 days) | Forever, and the group empties | With no DLQ, until retention (4 days); the mates behind it are never handled |
| What the default metric shows | Age metric: blind to it after 3 receives | Age: sees it | Lag: sees it (while the group settles) | Age: blind to its batch after 3 receives (a standard queue), then DLQ depth |
| Where it lands when bounded | The DLQ, alone | The DLQ, breaking the group's order | Your dead_letters table | The DLQ, with its mates unless reported by ID |
| Does order survive? | No order to keep | Only if you park | Only if you park | Only if you stop at the first failure (FIFO) |
| Replay | Redrive a new message | Redrive, out of order | Seek back to an offset | Redrive |
| What you must build | A DLQ alarm, retention, a runbook | Park or block, plus alarms | Bounded retries, the boundary, lag alarms | Per-record catch and IDs, concurrency caps |
Each column has its strength: the log wins replay and cheap fan-out, the standard queue wins "only this message waits", and FIFO gives order per group without a boundary you build.
Kinesis and DynamoDB Streams through Lambda retry a failing batch until its records expire by default ("a bad record can block processing on the affected shard for up to one week", AWS); bisecting, the maximum record age and the on-failure destination are Change Streams & the Transactional Outbox's (Part 9 there).
Your consumer retries a failing record in a loop. Why does a third of your traffic stop at once, and why does the group lose a consumer every five minutes until nothing reads at all?
What to remember from Part 8
- On a log, one bad record stops every partition its consumer owns, not only its key.
- Commit after the effect; auto-commit is safe only if everything polled is finished before the next poll.
- Bound in-line retries so
poll()keeps coming; then dead-letter, block the key, park.
Part 9. The backlog behind it
Replay B runs on a copy of Q4 without o-7731: from 20:00:00 to 20:10:00 a flash sale sends 100 orders a second against a fleet that finishes about 79, then traffic drops back to 15. At 20:07 about 9,000 orders are waiting (9,085 in seed 1, 8,621 in seed 2; 21 a second × 420 s ≈ 8,800 by arithmetic). How late are they, and which number should page someone?
Two different ages
A queue's backlog answers two questions, and they have different answers.
Formula 1: how long a message arriving now will wait. It sits behind the whole backlog, which drains at the fleet's throughput:
Formula 2: how old the oldest message is. The backlog was built by arrivals, so its oldest message arrived about one backlog's worth of arrivals ago:
They agree only when arrivals ≈ throughput. At 20:10 the backlog is about (100 − 79) × 600 = 12,600. A new order waits about 12,600 ÷ 79 ≈ 160 s; the oldest order is about 12,600 ÷ 100 ≈ 126 s old. And after the burst the oldest keeps ageing: the fleet now works through messages that arrived at 100 a second, advancing 79 ÷ 100 of a second of arrivals per second, so the oldest ages by 0.21 s every second until the fleet reaches messages that arrived after 20:10, about 160 s later. It peaks near 126 + 0.21 × 160 ≈ 160 s at about 20:12:40.
| # | Time | Event (seed 1; seed 2) |
|---|---|---|
| 30 | 20:00:00 → 20:10:00 | 100 orders a second against a capacity of about 79 (the run: 78.7 and 79.4 finished a second); 15 a second after 20:10 |
| 31 | 20:07:00 | Backlog 9,085; 8,621 |
| 31 | 20:09:08; 20:09:36 | The oldest message passes 120 s (arithmetic: 0.21 × t > 120 when t ≈ 571 s, 20:09:31). The age alarm (Q4: > 120 s) fires at the 20:10:00 datapoint in both seeds |
| 31 | 20:10:00 | Backlog 13,118; 12,345. Oldest 130 s; 125 s |
| 31 | 20:12:47; 20:12:38 | The oldest peaks at 168 s; 158 s, 2.6 to 2.8 minutes after the burst ended. The SLO (5 minutes) holds |
| 32 | 20:13:25; 20:13:14 | Empty. It drained at 79 − 15 = 64 a second: 12,600 ÷ 64 ≈ 197 s, so about 20:13:17 by arithmetic. The order enqueued at 20:10:00 waited 167 s; 157 s for its first receive (Formula 1: about 160 s) |
Synthesizing vector architecture diagram...
What to notice: the line that climbs to 166 at second 760 is replay B as described. It keeps rising for more than two minutes after the burst ends at second 600, then falls to zero as the fleet reaches the orders that arrived after the burst. The line that turns down at about 80 is the same burst with scaling on backlog per worker (side row 9s, below). The alarm line is 120 s and the SLO is 300 s; neither is drawn.
Drain time, in words: a backlog drains at capacity minus arrivals, not at capacity. Change Streams & the Transactional Outbox states it for a relay (Part 3 there) and Retries, Timeouts, Backpressure & Load Shedding for request queues (Part 6 there). The drill The Warehouse System That Fell a Day Behind applies it: its backlog grows by 2,200 − 400 = 1,800 orders a minute during a peak of T minutes and shrinks by only 400 − 150 = 250 a minute afterwards, so it drains in 1,800 T ÷ 250 = 7.2 T.
Age is the SLO; depth isn't
Depth says how much work is waiting; age says how late it is, in the units of the SLO. The same 12,600 messages clear in about 160 s at 79 a second whatever built them, but their oldest is 126 s old after a 100-a-second burst, and 12,600 ÷ 15 = 840 s, 14 minutes, if the fleet had stopped and they built up at 15 a second. Alarm on age. But know what each signal can't see:
| Signal | Units | Blind spot |
|---|---|---|
| Age of the oldest message, standard queue | Seconds | A message received 3 or more times is left out: the poison message (Part 5) |
| Age of the oldest message, FIFO queue | Seconds | Nothing like that: a blocked group raises it (Part 7) |
Depth (ApproximateNumberOfMessagesVisible) | Messages | Not in SLO units; the same depth can be 2 minutes or 14 minutes of lateness |
In flight (ApproximateNumberOfMessagesNotVisible) | Messages | Messages in backoff count too; a busy fleet and a stuck one look alike |
Consumer lag in records (MSK MaxOffsetLag) | Records | Emitted only while the group is STABLE or EMPTY; missing for a group that never settles (Part 8) |
Consumer lag in time (MSK EstimatedMaxTimeLag) | Seconds | The same gap |
| Every CloudWatch alarm | 1-minute datapoints | A breach inside a minute is seen at the minute's end: event 31's 120 s was crossed at 20:09:08 and alarmed at 20:10:00 |
For Kinesis Data Streams the matching signal is iterator age, and Event Time, Watermarks & Checkpoints explains how it differs from watermark lag (Part 12 there).
Q4 is the alarm set for the main line: age of the oldest > 120 s, DLQ depth ≥ 1, LastAttemptFailed ≥ 1, and more than 30 of the 40 slots in flight for 5 minutes. In replay B the in-flight alarm fires at 20:05:00, five minutes before the age alarm: a full fleet is the early warning, a late message the SLO signal.
Adding consumers
Side row 9s. AWS's method for scaling a queue's consumers (EC2 Auto Scaling) is backlog per worker: the acceptable backlog per worker = the latency you accept ÷ the time one worker takes per message. Here, 60 s of work per worker: a worker finishes 10 ÷ 0.508 s ≈ 19.7 orders a second, so 60 × 19.7 ≈ 1,181 messages per worker. Each minute the controller asks for backlog ÷ 1,181 workers, and a new worker is ready 2 minutes after it is asked for (our figure).
| Replay B, seeds 1; 2 | Without scaling | With scaling on backlog per worker |
|---|---|---|
| Workers added | 0 | 3, ready at 20:06, 20:07 and 20:08 (7 in all) |
| Backlog at 20:10 | 13,118; 12,345 | 2,437; 1,677 |
| Oldest age, peak | 168 s; 158 s at about 20:12:40 | 81 s; 78 s, at 20:06:59 and 20:06:34 |
| Age alarm (120 s) | Fires 20:10:00 | Never fires |
| Empty at | 20:13:25; 20:13:14 | 20:10:21; 20:10:14 |
Scaling lags a burst by the launch time plus a datapoint or two, so it trims the peak rather than preventing it. A Lambda consumer ramps by up to 300 concurrent invokes a minute (Part 6), and a log's consumer group can't use more consumers than it has partitions (Part 8).
The same 12,600-message backlog: how old is its oldest message after a 100-a-second burst, and after a slow build-up at 15 a second? And how long will a new order wait in each case?
What to remember from Part 9
- Alarm on how long the oldest work has waited, in the units of your SLO.
- The oldest message's age isn't backlog ÷ throughput, and it keeps rising after the burst.
- A backlog drains at capacity minus arrivals, not at capacity.
Part 10. The copies it fans out to
At 19:59 a deploy replaces the email queue's access policy, and the new policy drops the statement that lets the orders topic send to it. From 20:00 to 20:20, about 18,000 order confirmations are never sent, o-7731's among them. The publish succeeded, fulfil got every order, and nothing retried.
One publish, one copy per queue
SNS fan-out delivers one copy of each published message to every subscription: here three SQS queues. Each copy then lives its own life: its own visibility timeout, its own receive count, its own DLQ. A copy is wrapped in an SNS JSON envelope unless the subscription turns on raw message delivery. A filter policy on a subscription (on message attributes, or on the message body) decides which copies it gets: b2b takes only messages with channel = "b2b".
Synthesizing vector architecture diagram...
What to notice: the broken policy is on the email queue, not on the topic, so fulfil and b2b keep working. The dashed subscription DLQ exists only in Q5, and the drain script sends its messages straight to the email queue, never back to the topic.
What SNS and EventBridge retry, and what they drop
SNS sorts delivery failures into two kinds:
- Server-side errors (the endpoint's service is unavailable): to SQS and Lambda, SNS retries 3 times at once, 2 times 1 s apart, 10 times with backoff from 1 to 20 s, and 100,000 times 20 s apart: "100,015 times, over 23 days".
- Client-side errors, "when an owner deletes the endpoint … or when an owner changes the policy attached to the subscribed endpoint in a way that prevents Amazon SNS from delivering": SNS doesn't retry, and it discards the message "unless a dead-letter queue is attached to the subscription".
A queue encrypted with a customer managed KMS key is similar: SNS can send to it only if the key's policy grants the SNS service principal GenerateDataKey and Decrypt. AWS doesn't say whether a refused key is retried, which is one more reason to give every subscription a DLQ.
| # | Time | Setting | Event (seeds 1; 2) |
|---|---|---|---|
| 33 | 20:00:00 → 20:20:00 | Q0 | Client-side error on every email delivery: 18,228; 17,958 confirmations dropped (15 a second × 1,200 s = 18,000). No retry, no error in the publisher, fulfil unaffected |
| 34 | 20:00:00 → 20:28 | Q5: a DLQ on every subscription | The same 18,228 copies land in email-sub-dlq; its alarm fires at 20:01:00 (or 20:16:00 if CloudWatch takes its full 15 minutes for an idle queue). The policy is fixed at 20:20. At 20:25 a drain script moves the copies, oldest first, at 100 a second, into the email queue; it finishes at 20:28:02. The emails go out 8 to 25 minutes late; none is lost |
Why a drain script, not a redrive? SQS's redrive task (StartMessageMoveTask) accepts only DLQs "whose sources are other Amazon SQS queues"; a subscription DLQ's source is an SNS subscription. And never re-publish the messages to the topic: fulfil and b2b would get second copies.
textDRAIN A SUBSCRIPTION DLQ (email-sub-dlq -> email) after the cause is fixed, and after one message has gone through end to end: loop until email-sub-dlq is empty: msgs = ReceiveMessage(email-sub-dlq, up to 10, long poll) for each msg: SendMessage(email, msg.body) the original body; raw delivery keeps it unwrapped DeleteMessage(email-sub-dlq, msg) only after the send succeeded pace to 100 a second, below what the email senders can take
| On equal terms | SNS to an SQS or Lambda subscriber | EventBridge rule target (a classic event bus) |
|---|---|---|
| Retries after a service-side failure | 100,015 attempts over 23 days | Up to 24 hours and 185 attempts, "with an exponential back off and jitter" |
| After the retries run out | Discarded unless the subscription has a DLQ | "The event is dropped" unless the target has a DLQ |
| A permission error or a deleted endpoint | Not retried: DLQ or discarded | Not retried: sent "directly to the target DLQ", or dropped |
| Where failed messages land | A DLQ per subscription | A standard SQS DLQ per target |
| Bringing them back | Your own drain script | Your own script: the DLQ's source isn't an SQS queue either |
| Size limit | 256 KiB per message | A PutEvents request must be "less than 1 MB" |
| Filtering | A filter policy per subscription, on attributes or the body | An event pattern per rule, up to 5 targets per rule |
Those EventBridge figures describe the default bus and custom buses of the Classic type. EventBridge now also offers a newer Custom Event Bus with retention, FIFO delivery within an event group and 5-minute publish-time deduplication; its subscribers retry 5 times or for 300 s by default (settable up to 185 attempts and 24 hours) and drop an event that runs out, leaving only a metric, unless you give them a DLQ.
Size limits
| Hop | Largest message |
|---|---|
| SQS (standard and FIFO) | 1 MiB (1,048,576 bytes), raised from 256 KiB on 4 August 2025. (AWS's General Reference quota page still says 256 KB; the SQS guide and the announcement say 1 MiB) |
| SNS | 256 KiB |
EventBridge (classic PutEvents) | A request under 1 MB |
| Kinesis Data Streams | A record of up to 10 MiB |
| Kafka | About 1 MB by default (the broker's message.max.bytes, 1,048,588) |
| Lambda | 6 MB for a synchronous invoke's payload |
A message that fits the queue can fail at the topic in front of it: SQS takes 1 MiB, but SNS stops at 256 KiB.
Side row 10c. At 20:05:00 a B2B order, o-7790, is 1.4 MB: over SNS's 256 KiB and SQS's 1 MiB. The outbox's size check (page 05) rejects it at the insert, so it never becomes a message that can't be sent. The fix is a claim check: store the payload in S3 and send a pointer of about 0.4 KB.
textCLAIM CHECK producer: PUT the payload to S3 (bucket, key, version); compute a checksum publish {order, bucket, key, version, checksum} about 0.4 KB consumer: GET the object by bucket, key and version; check the checksum; process never delete the object: every copy of the pointer must still find it S3 lifecycle rule: delete objects after 30 days
Why 30 days? The object must outlive every copy of the pointer. A copy that fails lands in a standard DLQ, where it can live 14 days counted from its original enqueue time; after one redrive it is a new message, and if it fails again it can live up to 14 more days in the DLQ: 14 + 14 = 28 days, so a 30-day rule. With fan-out, several consumers read the same object, so no consumer may delete it. The SQS Extended Client Library for Java deletes the S3 object when you delete the message, by default: turn that off. (The Python client's default is off.) Step 1.3 of the email loop keeps raw mail in S3 and passes a reference for the same reason.
The publish succeeded and fulfil got every order. Why did 18,000 customers get no confirmation, why did nothing retry, and why can't you just redrive them?
What to remember from Part 10
- Fan-out makes one copy per subscriber; each copy can fail on its own.
- Give every subscription a DLQ and a drain script: some errors are never retried.
- Check every hop's size limit; above it, send a pointer to S3.
Part 11. One order, end to end
With Q5, every hop has a guard. Here is o-7731 through all of them (seed 1), from the checkout to the pick list.
Synthesizing vector architecture diagram...
What to notice: every arrow can repeat, and the boxes below each arrow are what make the repeat harmless. The one green box, the delete, comes last, after the effect is durable.
| Time (seed 1) | What happens to o-7731 |
|---|---|
| 20:00:07.05 | The checkout commits the order and its outbox row in one transaction |
| 20:00:07.65 | The relay publishes OrderPlaced to orders |
| 20:00:07.70 | SNS delivers a copy to fulfil and one to email; a worker receives the fulfil copy at once |
| 20:00:08.10 | The WMS creates PL-501, then hangs; the worker times out after 5 s |
| 20:00:07.70 → 20:01:52.41 | Five deliveries, each failing in 5 s, backed off with jitter |
| 20:02:00 | LastAttemptFailed alarm; the fulfil on-call pages the WMS team |
| 20:02:16.60 | Moved to fulfil-dlq (14-day retention); the DLQ alarm fires at 20:03:00 |
| 21:30:00 | The WMS fix is deployed |
| 21:40:00 | Redriven at 10 a second: a new message, receive count 0, received at once |
| 21:40:00.52 | The WMS returns PL-501 by its key, the label prints, the pick and the inbox row are written |
| 21:40:00.53 | Deleted: the ack point |
On the main line, none of Q4's backlog alarms fires: at 15 a second the fleet never runs short, and o-7731 is out of the age metric from its 3rd receive. The alarms that caught it were the two built for a poison message.
| Hop | Ack point | What can repeat | The guard | Lost without the guard |
|---|---|---|---|---|
| Checkout → outbox | The commit | Nothing: one transaction | One transaction for the order and its event (Change Streams) | An order with no event, or an event with no order |
| Relay → SNS | Marking the outbox row sent, after SNS accepts | A publish, if the relay dies before marking | The same event ID on every re-publish | Nothing lost; duplicates downstream |
| SNS → each queue | SNS's delivery | A copy may be delivered twice (at-least-once) | A DLQ on every subscription; the inbox downstream | Every copy that hits a client-side error |
fulfil → worker | DeleteMessage, after the effect | Redelivery after a lease runs out; a crash after the effect | Visibility + heartbeat, the inbox check, maxReceiveCount, a watched 14-day DLQ | A poison message loops until retention deletes it |
| Worker → WMS | The WMS's response | A call that timed out (an unknown) | Idempotency-Key: o-7731:pick | A second pick list |
| Worker → fulfil DB | The transaction | Nothing: the inbox row is in it | The inbox row, unique per order | A second pick recorded |
email → email sender | DeleteMessage after the send | A send, if the worker dies after sending | The provider can't take a key: a rare duplicate email is accepted (Idempotency, Part 7 there) | A missing confirmation |
What to remember from Part 11
- Say delivery semantics hop by hop.
- Every hop repeats; the effect decides whether a repeat hurts.
- The ack point is the last step, after the effect is durable.
Part 12. On AWS
Every mechanism on this page is a managed feature somewhere in AWS. Three rules carry over: SQS, SNS and EventBridge deliver at least once, so you still choose the ack point and make the effect safe to repeat; each service has its own place where failed work lands, and each must be set and alarmed; and their limits differ, so check each hop.
Managed services that use it
| Service | What it provides | What AWS documents |
|---|---|---|
| Amazon SQS (standard) | Leased delivery, redrive to a DLQ, delays, long polling | At-least-once: "more than one copy of a message might be delivered, and messages may occasionally arrive out of order"; on rare occasions a message can be received after it was deleted. Messages are stored in several Availability Zones before a send is acknowledged. Limits in the table below |
| Amazon SQS (FIFO) | Order per message group; deduplication of sends | While a group has a message in flight "subsequent messages in that group are not made available"; a receive "attempts to return as many messages as possible with the same message group ID"; deduplication IDs are remembered for 5 minutes, even after the message is deleted; no per-message timers; 300 transactions a second per API action (3,000 messages with batches of 10), high-throughput mode up to 70,000 in three Regions; a FIFO queue needs a FIFO DLQ, and "don't use a dead-letter queue with a FIFO queue if you don't want to break the exact order" |
| SQS dead-letter queues and redrive | Setting aside and bringing back | maxReceiveCount is "the number of times a consumer can receive a message from a source queue before it is moved"; the move happens at a later receive, and viewing messages in the console counts as receiving them; same account and Region; a standard DLQ keeps the original enqueue time, a FIFO DLQ resets it; AWS advises a DLQ retention longer than the source's; StartMessageMoveTask only for DLQs whose sources are SQS queues, up to 500 messages a second, at most 36 hours per task, up to 100 active tasks per account; redriven messages are new messages |
| SQS in CloudWatch | The signals | ApproximateAgeOfOldestMessage: on standard queues, messages received 3 or more times are moved to the back and excluded until processed; not on FIFO; on a DLQ, measured from the move. NumberOfMessagesSent doesn't count redrive-policy moves. 1-minute metrics; a queue inactive for 6 hours stops reporting, and its first datapoint after it becomes active can be up to 15 minutes late |
| SQS fair queues | Tenant fairness on a standard queue | MessageGroupId = tenant; SQS favours quieter tenants when one holds a disproportionate share of in-flight messages; no order and no per-tenant rate limit |
| Amazon SNS | Fan-out, filtering, delivery retries | One copy per subscription, in a JSON envelope or raw; filter policies on attributes or the body; 256 KiB messages (the Extended Client Library can carry up to 2 GB through S3); to SQS and Lambda, 100,015 attempts over 23 days for server-side errors; client-side errors not retried and discarded unless the subscription has a DLQ; an encrypted queue's KMS key must grant SNS GenerateDataKey and Decrypt. FIFO topics also exist, with their own throughput limits |
| Amazon EventBridge | Rules that route events to targets | On a classic event bus: retries "for 24 hours and up to 185 times", then "the event is dropped" unless the target has a DLQ (a standard SQS queue); permission errors and missing targets skip retries and go "directly to the target DLQ"; a PutEvents request under 1 MB; up to 5 targets per rule. The newer Custom Event Bus adds retention, FIFO delivery per event group and 5-minute deduplication, with shorter default retries (5 attempts or 300 s) |
| AWS Lambda (event source mappings) | A managed consumer for SQS, Kinesis, DynamoDB Streams, Amazon MSK and Amazon MQ | For SQS: batches up to 10,000 (10 for FIFO), a window up to 5 minutes (standard queues only); visibility at least 6 × the function timeout + the window; maxReceiveCount at least 5; 5 concurrent batches at first, up to 300 more a minute, up to 1,250; maximum concurrency 2 to 1,000, not above reserved concurrency; partial batch responses and their success and failure rules; on FIFO, concurrency ≤ active groups and "stop processing messages after the first failure". Provisioned mode: dedicated pollers, each up to 1 MB a second and 10 concurrent invokes. On-failure destinations exist for Kinesis, DynamoDB Streams and Kafka sources, not SQS |
| Amazon MSK | Managed Kafka: consumer groups and offsets | Consumer-lag metrics (MaxOffsetLag, SumOffsetLag, OffsetLag, EstimatedMaxTimeLag, EstimatedTimeLag), "emitted only if a consumer group is in a STABLE or EMPTY state". Kafka's client defaults apply: max.poll.interval.ms 300,000, session.timeout.ms 45,000, auto-commit every 5,000 ms, max.poll.records 500, Range first in the assignor list |
| Amazon Kinesis Data Streams | Ordered shards read by position | Records up to 10 MiB; retention 24 hours to 365 days; 2 MB a second of reads per shard, shared by all standard consumers; with Lambda, a failing batch retries until its records expire unless you bisect or cap retries (Change Streams, Part 9 there) |
| Amazon MQ (for RabbitMQ) | Managed queues with acknowledgements | RabbitMQ 4.3, 4.2 and 3.13. consumer_timeout defaults to 30 minutes: a consumer that doesn't acknowledge in time has its consumer cancelled (or its channel closed) and its messages returned. Quorum queues take a delivery-limit policy; over the limit a message is dropped, or dead-lettered if a dead-letter exchange is set (RabbitMQ 4.x defaults the limit to 20). Cluster brokers get a max-length policy with overflow: reject-publish, so a full queue rejects new messages |
| Amazon S3 | The claim check's payload store | The SQS Extended Client Library (Java and Python) puts payloads up to 2 GB in S3 and sends a reference; the Java client deletes the object when the message is deleted unless you turn that off |
SQS limits in one place
| Limit | Value |
|---|---|
| Visibility timeout | 0 s to 12 hours, default 30 s; "12 hours from when the message is first received", and extending doesn't reset it |
| Retention | 60 s to 14 days, default 4 days |
| Delay | 0 to 15 minutes, per queue or (standard only) per message |
| Long polling | Up to 20 s |
| Messages per receive, per batch call | 1 to 10 |
| Message size | 1 MiB |
| Backlog | "Unlimited" |
| In flight | About 120,000 (standard and FIFO) |
| A delete with an old receipt handle | "The request will succeed, but the message might not be deleted" |
Running it yourself
| Option | What it is | Facts and sizing |
|---|---|---|
| Kafka on Amazon EC2 or Amazon EKS | Self-run brokers, consumer groups and offsets | The same client settings as MSK. You run the brokers (KRaft), their disks (EBS or instance-store NVMe, see Write-Ahead Log, fsync & Group Commit) and the lag metrics yourself (Burrow, or an exporter). Create at least as many partitions as the largest consumer group you plan, and leave headroom for replays |
| RabbitMQ on EC2 or EKS | Self-run queues with acknowledgements and prefetch | Quorum queues replicate with Raft. The delivery limit defaults to 20 on 4.x; since 4.3 a basic.nack doesn't count toward it (a basic.reject does), so a consumer that nacks can loop without end. Over the limit a message is dropped unless a dead-letter exchange is configured. consumer_timeout 30 minutes; prefetch bounds what each consumer holds; max_message_size defaults to 16 MiB |
| Sizing in words | Workers from Little's law (arrival rate × time per message ÷ slots per worker), plus headroom for the drain you want (capacity − arrivals); visibility above the slowest normal job, with a heartbeat; DLQ retention 14 days; alarms on age, DLQ depth and your own last-attempt metric; scale on backlog per worker |
Queues and their DLQs are Regional: a Region failover needs its own plan for the messages left behind (Multi-Region Failover).
Look-alikes that are not this mechanism
| Look-alike | Why it looks like this | Why it isn't |
|---|---|---|
| AWS Step Functions "redrive" | The same word | Restarts a failed Standard workflow execution from the failed state; it doesn't move messages out of a DLQ |
| Lambda's asynchronous event queue, a function's DLQ, on-failure destinations | "Lambda has a DLQ" | They apply to asynchronous invokes (SNS → Lambda, for example). For an SQS source, failures land in the queue's redrive DLQ; event source mapping destinations exist only for Kinesis, DynamoDB Streams and Kafka |
| Amazon EventBridge Scheduler | "A delayed message" | A scheduler for one-time or recurring API calls, which SQS's guide suggests for delays beyond 15 minutes; it has no lease or receive count |
| "SQS FIFO exactly-once processing" | The name | Deduplication of sends for 5 minutes; consumers still see redeliveries (Idempotency) |
| Amazon Data Firehose | "It buffers and delivers records" | A delivery pipeline into stores; there are no consumers, leases or DLQ semantics to design |
| DynamoDB Streams | "A queue of changes" | A change log with a fixed 24-hour retention, read by position (Change Streams) |
| S3 Event Notifications | "S3 sends a message" | At-least-once notifications into SQS, SNS, Lambda or EventBridge; the queue behind them is the mechanism |
| Amazon ElastiCache (Valkey) lists or streams | "Valkey as a queue" | You build the lease (reclaiming pending entries), the receive count and the DLQ yourself; there is no managed redrive |
What to remember from Part 12
- SQS, SNS and EventBridge are at-least-once; you choose the ack point and the guard.
- Their limits differ: visibility 12 h, delay 15 min, retention 14 days, size 1 MiB vs 256 KiB.
- Each service has its own place where failed work lands: set every one and alarm on it.
Part 13. What you've learned
Back to o-7731
As shipped, o-7731 failed 11,520 times over four days, held 1.33 worker slots and WMS connections the whole time, kept every alarm green and was deleted by retention on Tuesday. With every guard in place (Q5), it was set aside within minutes, someone was paged at 20:02, the WMS team fixed the bug at 21:30, and a redrive picked it at 21:40:00.52 with the one pick list the WMS had created at 20:00:08. Here is what each piece did:
- The ack point after the effect (Part 1) turned a crash from a lost order into a 30 s delay.
- Parallel batches, or receiving one at a time (Part 2), stopped healthy messages from outliving their leases behind a slow one.
- A sized lease with a heartbeat, and the inbox check first (Part 3) rendered
o-7740once, and kept every repeat from making a second pick. - Counting, backoff and
maxReceiveCount(Part 4) cut 11,520 WMS attempts to 5; deferring by re-sending or pausing kept 5,400 healthy orders out of the DLQ during a throttle. - A watched, 14-day DLQ with a last-attempt metric (Part 5) paged someone at 20:02, and a capped redrive after the fix finished the job.
- Partial batch responses (Part 6), parking a FIFO group (Part 7) and a bounded error boundary on a log (Part 8) kept one bad message from taking its batch, its customer or a third of all customers with it.
- Age alarms and scaling on backlog per worker (Part 9) watched lateness in SLO units.
- A DLQ and a drain script on every subscription, and a claim check (Part 10) kept fan-out from dropping anything silently.
The snapshots, side by side
| Snapshot | Setting and time | What it showed |
|---|---|---|
| M1 | Q0, 20:00:08.70 | Received once; PL-501 created; the WMS call hanging |
| M2 | Q0, 20:00:52.70 | Two workers have held it; receive count 2; W4 still hanging |
| M3 | Q0, Saturday 20:00 | Receive count 2,880; 1.33 slots and connections held; the age alarm green because the message is excluded |
| M4 | Q2, 20:02:16.60 | In the DLQ after 5 deliveries; PL-501 unlabelled |
| M5 | Q3, 21:40:01 | Redriven as a new message, picked, one pick list, DLQ empty |
| M6 | F1 vs F3, 20:02:30 | c-19 blocked with its later messages waiting, vs parked in order while every other customer flows |
What it costs
- Duplicates to make harmless: every consumer needs an inbox or a key, because at-least-once means some messages arrive twice.
- A visibility timeout to tune: too short gives duplicates and stale deletes; too long makes a crashed worker's message wait. A heartbeat costs a thread and an API call per message.
- An owner for every DLQ: an alarm, a runbook, and someone who can fix the cause and redrive.
- Healthy work waits during deferral: pausing or re-sending keeps orders out of the DLQ but makes them late (10 to 11 minutes in replay T).
- Ordered queues and partitions block: a failing message stops its group or its consumer's partitions, and parking needs tables and a repair job.
- Metrics that lie in documented ways: the age metric's poison exclusion,
NumberOfMessagesSentignoring moves, late datapoints from idle queues, missing lag for groups that never settle, and an age that resets when a message is re-sent.
The whole story, event by event
Side rows and side replays run on copies and are marked "(side)". Seeds 1 and 2 unless stated; times are seed 1's.
| # | Time | Setting | Event |
|---|---|---|---|
| 1 | about 20:00:07.05 | Q0 | o-7731's checkout commits with its outbox row |
| 2 | 20:00:07.70 | Q0 | Relay → SNS → fulfil and email, not b2b (its filter) |
| 3 | 20:00:07.70 (t1) | Q0 | Received, count 1, hidden until 20:00:37.70 |
| 1x | (side) | Delete on receive | The worker holding o-7733 is killed 0.2 s after the receive: o-7733 is never picked, and the worker's other 1 or 2 orders are lost too |
| 1y | (side) | Delete after | The same kill: o-7733 back at 20:00:50.65, picked at 20:00:51.12; one pick list |
| 4 | 19:58 → 20:10 | Q0 | Every receive returns 1 message; about 7.6 of 40 slots busy; capacity about 79 a second |
| 2b | (side) | Q0, forced batch | W2's batch of 10 at 20:00:08, worked serially: the 5 behind o-7731 are received again at 20:00:38 and picked by others 30 s late; W2's later deletes are no-ops |
| 2s | (side) | Night | Short polling 2,400 receive calls a minute for one order; long polling about 13 |
| 5 | 20:00:08.10 | Q0 | The WMS creates PL-501, then hangs |
| 6 | 20:00:37.70 | Q0 | Received again (count 2): two workers hang on one order |
| 7 | 20:00:47.70 | Q0 | The first worker's 40 s timeout: no delete |
| 8 | 20:01:07.70 | Q0 | Count 3: out of the age metric |
| 9 | from then on | Q0 | 1.33 slots and WMS connections held; 2,880 deliveries a day |
| 10 | 20:00:20.65 → 20:01:50.66 | Q0 | o-7740: 3 renders, 4 deliveries; the first pick's delete (20:01:36.27) is a no-op; the inbox check deletes it at 20:01:50.66; one pick list |
| 11a | (side) | Q1, no heartbeat | o-7740: 2 renders, deleted 20:02:20.66 |
| 11 | 20:01:36.27 | Q1 | o-7740 with a heartbeat: 1 render, deleted with a current handle. o-7731: 0.08 slots, 1,440 deliveries a day |
| 3z | (side) | text | A job longer than 12 hours: claim it, delete the message, keep your own lease |
| 12 | 20:00:07.70 → 20:02:16.60 | Q2 | 5 deliveries with jittered backoff; moved to fulfil-dlq at t1 + 128.9 s (seed 2: 201.6 s). Over 5,000 seeds: mean 179.7 s, 38.7 to 318.6 s |
| 13 | Q0 vs Q2 | WMS label attempts: 11,520 vs 5 | |
| 4t | (side) | Replay T | A 10-minute 429 throttle: 5,374 and 5,421 healthy orders dead-lettered when deferring by visibility (7,904 and 8,017 with jitter); 0 when re-sending with a delay (longest wait about 655 s) or pausing (601 s) |
| 4d | (side) | text | DelaySeconds ≤ 15 minutes; no per-message timers on FIFO; longer waits belong to a scheduler |
| 14 | Friday → Tuesday | Q0 | 11,520 deliveries; the age alarm green; deleted by retention on Tuesday at 20:00:07.70 |
| 5a | (side) | Q0 + a DLQ on Saturday | Moved on Saturday at 20:00:07.70; the NumberOfMessagesSent alarm never fires; the DLQ's age reads 72 h when it expires on Tuesday at 20:00:07.70 |
| 15 | 20:01:52.41 → 20:03:00 | Q3 | LastAttemptFailed alarm at 20:02:00; the move at 20:02:16.60; the DLQ-depth alarm at 20:03:00 (or 20:18:00 with the inactive-queue delay) |
| 16 | 21:30:00 | Q3 | Paged, the WMS team deploys the fix |
| 17 | 21:40:00 → 21:40:00.53 | Q3 | Redrive at 10 a second: a new message, picked with PL-501 at 21:40:00.52; 1 h 40 min late, one pick list |
| 5c | (side) | Q3, redrive at 21:00 | Five more hangs; back in the DLQ at 21:04:05 (seed 2: 21:03:25) |
| 5d | (side) | text | Redrive: up to 500 a second, at most 36 hours a task, new messages that interleave with new traffic, SQS-sourced DLQs only |
| 5n | (side) | Q3, consumers stopped at 20:02 | o-7731 is never moved and the DLQ alarm stays silent; the age alarm (> 300 s) fires at 20:08:00 only because healthy orders pile up |
| 18 | 20:00:07 | Replay L | Lambda: batch 10, window 1 s, function timeout 15 s, visibility 100 s; o-7731 5th in its batch |
| 19 | 20:00:08 → 20:08:28 | L1: throws | The batch re-forms 5 times: m1 to m4 handled 5 times (16 repeat runs), m6 to m10 never; all 10 dead-lettered at 20:08:28. Averages over positions: 18, or 13.8 shuffled |
| 20 | 20:00:18.1 | L2: returns [] | All 10 deleted: o-7731 lost, not in the DLQ |
| 21 | 20:00:08 → 20:08:28 | L3: returns its ID | 9 mates deleted after one pass; o-7731 alone to the DLQ; 0 repeats |
| 6r | (side) | text | Throttling by reserved concurrency sent 9 healthy messages to a DLQ in AWS's demonstration; maximum concurrency sent none |
| 22 | 20:00:07.70 on | Replay F | fulfil.fifo, group = customer; c-19's group blocked, every other group flows |
| 23 | Friday → Tuesday | F1: no DLQ | o-7802 picked on Tuesday at 20:00:41.36, four days late; a 120 s age alarm fires at 20:03:00 |
| 24 | 20:02:16.60 → 21:40 | F2: FIFO DLQ | The slot change fails 5 times and follows o-7731 into the DLQ (20:05:48.28); o-7802 picked at 20:05:48.68, ahead of the order before it |
| 25 | 20:01:52.41 → 21:40 | F3: block and park | o-7731 dead-lettered in the consumer's table; the slot change and o-7802 parked in order; the repair applies all three in order |
| 7k | (side) | F1, other keys | By customer: 2 later messages wait 4 days. By order: only the slot change waits. One group: everything waits, and even without a poison message it manages about 2 a second |
| 7t | (side) | text | 300 transactions a second per API action; up to 70,000 in high-throughput mode; heartbeats spend the budget |
| 26 | 20:00:07 | Replay K | 6 partitions, 3 consumers with Range; auto-commit; max.poll.interval.ms 300 s |
| 27 | 20:00:08.48 | K1 | C2 stuck: partitions 2 and 3 stop; lag +5 a second |
| 28 | 20:05:08.48 → 20:15:25.22 | K1 | C2, then C1, then C3 declared failed; a half, then all partitions stop; the group is empty; lag +15 a second |
| 8a | (side) | K, a thread pool | Auto-commit past handed-off records: o-7731 and its slot change skipped for good when C2 is killed at 20:02:00 |
| 29 | 20:00:08.48 → 20:00:34.60 | K2 | 3 bounded retries (26.1 s), dead-letter, block c-19, park; partitions 2 and 3 wait once; no rebalance |
| 8k | (side) | text | Kinesis and DynamoDB Streams through Lambda block a shard until records expire unless bisected or capped |
| 30 | 20:00 → 20:10 | Replay B | 100 orders a second against about 79 |
| 31 | 20:07 → 20:12:47 | B | Backlog 9,085 at 20:07 and 13,118 at 20:10; oldest 130 s at 20:10, peak 168 s at 20:12:47; age alarm at 20:10:00 |
| 32 | 20:10 → 20:13:25 | B | Drains at about 64 a second; the order enqueued at 20:10:00 waited 167 s |
| 9s | (side) | B + scaling | 3 workers added; peak age 81 s; empty at 20:10:21; the age alarm never fires |
| 9d | (side) | text | 12,600 messages: 126 s old after a 100-a-second burst, 14 minutes old if built at 15 a second |
| 33 | 20:00 → 20:20 | Q0 | The email queue's policy refuses the topic: 18,228 confirmations dropped, no retry |
| 34 | 20:00 → 20:28:02 | Q5 | Subscription DLQ, alarm at 20:01:00, drain script from 20:25: emails 8 to 25 minutes late, none lost |
| 10c | (side) | Page 05's relay | o-7790 (1.4 MB) over SNS's and SQS's limits: a claim check in S3, 30-day lifecycle, no consumer deletes it |
| 10e | (side) | text | A classic EventBridge target down for 30 hours: events older than 24 hours are dropped (185 attempts) unless it has a DLQ |
| 35 | 20:00:07.05 → 21:40:00.53 | Q5 | o-7731 through every hop, each with its ack point and guard |
The cheat card
| Topic | Remember |
|---|---|
| Semantics | Ack before the effect = at-most-once; after = at-least-once; + an effect that ignores repeats = effectively-once |
| Receive loop | Long poll 20 s; receive what you'll start; every timer starts at the receive; inbox first; delete per message after the effect |
| Visibility | Above the slowest normal job; heartbeat (60 s every 20 s here); 12 h cap from the first receive; release early with 0 |
| Stale handles | A delete with an old receipt handle "will succeed", and may delete nothing: guard the effect, not the delete |
| Slow failure | Slots held = attempt time ÷ visibility (40 ÷ 30 = 1.33); deliveries = retention ÷ visibility (4 days ÷ 30 s = 11,520) |
| Classify | Transient and throttle: retry later; permanent: DLQ now (send, then delete); unknown: retry with the key |
| Backoff | Visibility = random(0, min(900 s, 10 s × 2^(n−1))): about 180 s to the DLQ with 5 receives |
| Deferral | Every receive counts: defer healthy work by re-sending with a delay (≤ 15 min) or pausing, not by visibility |
| DLQ | maxReceiveCount ≥ 5; moved at the next receive; alarm on visible depth ≥ 1 and your own last-attempt metric; 14-day retention; an owner |
| Redrive | Look, fix, test one, capped rate, watch, note; redriven = new message; subscription DLQs need a script |
| Lambda | Visibility ≥ 6 × function timeout + window; function timeout covers a batch with one hang; report failed IDs; empty list = all succeeded; cap with maximum concurrency |
| FIFO | Order per group; group by the smallest thing whose order matters; a DLQ breaks order, parking keeps it; receive 1 per group |
| Log | One consumer per partition; one stuck record stops all its consumer's partitions; bound retries under max.poll.interval.ms; commit after the effect; cooperative assignor, static membership |
| Backlog | New arrival waits backlog ÷ throughput; the oldest is backlog ÷ arrival rate and keeps ageing; drain = backlog ÷ (capacity − arrivals) |
| Signals | Alarm on age; standard queues hide poison from it; MSK lag missing for unsettled groups; 1-minute datapoints |
| Fan-out | A DLQ per subscription; client-side errors never retried; drain into the queue, never back to the topic |
| Size | SQS 1 MiB, SNS 256 KiB, EventBridge < 1 MB, Kinesis 10 MiB, Kafka about 1 MB; above: a claim check in S3 |
Failure checklist
- Is every delete (or commit) after the effect is durable, and is every effect guarded by an inbox or a key?
- Does each worker receive only what it will start now, and work a batch in parallel?
- Is the visibility timeout above the slowest normal job, with a heartbeat that stops when the work does?
- Do you know that a delete with an old receipt handle may do nothing, and check the inbox before any work?
- Are errors classified, with permanent ones sent to the DLQ at once?
- Do real retries back off with jitter, and does deferral of healthy work avoid spending receives?
- Does every queue have a redrive policy with
maxReceiveCountof at least 5? - Is every DLQ alarmed on its visible depth (not
NumberOfMessagesSent), with your own last-attempt metric as well? - Is every DLQ's retention 14 days, longer than its source's?
- Does every DLQ have an owner, a runbook and a tested, rate-capped redrive, used only after the fix?
- Do Lambda handlers catch per record and return failed IDs, never an empty list after an error, with concurrency capped at the event source?
- On FIFO queues and log partitions, is there a decision for a failing message: block and page, or park the key, rather than a DLQ that silently breaks order?
- Are in-line retries on a log bounded well inside
max.poll.interval.ms, with commits after processing? - Does every SNS subscription have a DLQ, an alarm and a drain script, and does every hop's size limit fit the largest message?
Think-first drills
Drill 1. A queue has a 45 s visibility timeout and 7-day retention, and no DLQ. A message fails after 50 s of work every time. How many deliveries does it get, and how many workers does it hold on average? What changes with maxReceiveCount 4 and a 90 s visibility timeout?
Drill 2. A dependency throttles for 8 minutes. Each deferral sets the visibility to 120 s, maxReceiveCount is 3, and 20 messages a second arrive. How many healthy messages reach the DLQ, and what change sends none?
Drill 3. 1,500 messages a minute arrive for 20 minutes against a fleet that finishes 1,000 a minute; then 300 a minute arrive. How big is the backlog at the end of the burst, how old is its oldest message then and at its peak, and when is it empty?
Interview questions
| Question | Model answer |
|---|---|
| At-most-once, at-least-once, exactly-once: which does a queue give you, and where does the choice get made? | SQS, SNS and EventBridge give at-least-once. The consumer chooses by where it acknowledges: a delete (or offset commit) before the effect is at-most-once, after it is at-least-once. Exactly-once exists inside one store (a transaction); across a queue you get effectively-once by making the effect ignore repeats: an inbox row in the same transaction, or an idempotency key at the downstream API. FIFO deduplication only drops repeated sends for 5 minutes. |
| How do you choose a visibility timeout, and what do you do for a job that sometimes takes an hour? | Above the slowest normal processing time, not the average, and extended by a heartbeat while the work runs (ChangeMessageVisibility from a separate thread, stopped when the work ends or hangs). The heartbeat can extend up to 12 hours from the first receive. A job that can exceed that: claim it in your own table with your own lease, and delete the message. Either way, check the inbox before working, because a lease can still expire under a live worker and its delete may do nothing. |
| A message fails forever. Walk me from the first failure to it being processed after the fix, and name every alarm. | Classify the error; back off by setting the visibility with jitter; the receive count rises each time; with maxReceiveCount 5 it is moved to a DLQ at the next receive after the 5th. Alarms: my own "last attempt failed" metric (fires before the move and doesn't depend on the DLQ's metrics), and the DLQ's ApproximateNumberOfMessagesVisible ≥ 1 (never NumberOfMessagesSent, and an idle DLQ's first datapoint can be 15 minutes late). The DLQ keeps it 14 days; the on-call pages the owning team; after the fix, test one message, redrive at a capped rate and watch. The redriven message is new, and the idempotency key makes the retry safe. |
| A poison message on a FIFO queue, on a standard queue and on a Kafka partition: who else waits, and what do you do? | Standard queue: nobody else waits; it loops until retention unless a DLQ bounds it, and the age metric can't see it, so count receives. FIFO: its whole message group waits; a DLQ unblocks the group but breaks its order, so block and page, or park the group's later messages in your own table and repair in order. Kafka: every partition its consumer owns stops, and after max.poll.interval.ms the partition moves, with the record still at its head, to the next consumer: bound in-line retries, then dead-letter in your own database, block the key, park, commit. |
| Your Lambda consumer processes some orders twice, loses others, and puts healthy ones in the DLQ. What could explain each? | Twice: the handler throws, so the whole batch comes back and the records before the failure run again (or the visibility timeout is shorter than 6 × the function timeout + the window). Lost: the handler catches an error and returns an empty batchItemFailures list, which Lambda treats as complete success. Healthy ones in the DLQ: whole-batch failures raise every message's receive count, and throttling by reserved concurrency returns batches with higher counts. Fix: report failed IDs, size the timeouts, cap with maximum concurrency, maxReceiveCount ≥ 5, and an inbox. |
| Queue or log for this pipeline: fan-out, replay, ordering, one slow consumer? | A log if you need replay (seek back within retention), cheap fan-out to many teams (a consumer group each) and order per key at high throughput, and you'll build the error boundary. A queue (SNS in front for fan-out) if each message is an independent task, one slow or failing message must not stop others, and you don't need to replay. For per-key order on a queue, SQS FIFO groups. Either way the consumer is at-least-once and needs an idempotent effect. |
Where to go next
- Idempotency & Effectively-Once Processing: the keys and the inbox every redelivery relies on (Parts 1 and 6 there), and notifications that can't take a key (Part 7 there).
- Change Streams & the Transactional Outbox: the outbox that published
o-7731, the drain formula (Part 3 there), the error boundary on a log (Part 6 there) and Lambda on DynamoDB Streams (Part 9 there). - Leases, Fencing Tokens & Distributed Locks: the visibility timeout as a lease without a fence (Part 11 there).
- Retries, Timeouts, Backpressure & Load Shedding: error classes, full jitter and drain arithmetic (Parts 3, 4 and 6 there).
- Event Time, Watermarks & Checkpoints: iterator age vs watermark lag on a stream (Part 12 there).
- Sharding, Hot Keys & Rebalancing: a hot message group or partition key is a hot key.
- Multi-Region Failover: queues and their DLQs are Regional.
- Coming later: Rate Limiting Algorithms (the throttles that make you defer work in Part 4).
- Drill: The Warehouse System That Fell a Day Behind, both questions answered in Parts 1, 3, 4, 5 and 9.
- Background: Message Queues vs Event Streams, the broader primitive.
- Loops that rely on this page: the message queue (steps 1.2, 1.5, 1.6, 2.1 to 2.6, 3.6, R2.8, R2.9 and R3.9), notifications (steps 1.1 to 1.3, 2.1, 3.1, R1.9 and R3.9), the job scheduler (steps 1.3, 1.4, 2.5, 2.6 and 3.2), the web crawler (steps 1.1, 1.3, 2.3 and 2.5), email (steps 1.3 and 2.5), YouTube (step 2.2), the outbox ledger (steps 2.3, 3.5 and 3.6), ride-sharing (step 1.4), the news feed (step 1.4), chat (step 2.2), metrics and alerting (step 2.2), the URL shortener (steps 2.4 and 2.5), mobile chat, the digital wallet, the gaming leaderboard, hotel reservations, Google Drive, the offline-first news feed, the mobile stock-trading app, and the Uber and Netflix case studies.