Idempotency & Effectively-Once Processing
The Flight Booking That Charged Without a Seat
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
Tapped once, charged twice
A buyer taps Pay once for order 42, $30.00. The request reaches the load balancer 0.02 s later and then waits in a busy server's queue. At 10 s the phone gives up, asks "did it go through?", hears "not found", and the buyer taps again. That second request charges the card at 10.28 s. At 12.00 s the first request finally leaves the queue, and at 12.28 s it charges the card too. Tapped once, charged twice: $60.00.
Nobody wrote a bug here. Every part did something reasonable. Something repeated the request, and nothing could tell that the repeat was the same payment.
One tap, eight systems
The tap crosses eight systems before the money is settled, and then a ninth tells the buyer:
phone → load balancer → payment API → card processor → ledger → outbox relay → queue → worker, and the worker hands a push notification to a provider (Apple's or Google's push service).
The card processor is the PSP (payment service provider, for example Stripe or Adyen): the outside company that actually charges the card.
Every hop can lose an answer, and at every hop something may send the same thing again:
| Who may repeat | What it repeats |
|---|---|
| The phone's HTTP library | A request whose connection died |
| The app | A request that timed out, even days later |
| The API's resolver (a background job that settles unclear outcomes) | A charge whose reply was lost |
| The outbox relay | An event it published but didn't mark as done |
| The queue | A message whose worker died before deleting it |
| The worker | A push notification it may already have sent |
The number of unclear outcomes is not small. The payments loop assumes 50,000 payments a day and a 0.1% PSP timeout rate: 50,000 × 0.001 = 50 payments a day whose outcome we don't know, every day.
The phone sent Pay, heard nothing, and will send it again, maybe in 10 seconds, maybe in two days. Which piece of the system decides that the second send is the same payment, and how long must it remember?
The big picture
Synthesizing vector architecture diagram...
What to notice: one ID chain runs through the whole path. The phone's key K1 names the payment pay_7; pay_7 names the PSP charge and the event ev_7; ev_7 is what the queue, the worker's inbox and the notification record recognise. The ledger row is the durable guard: it exists for as long as the payment does.
Everyone remembers for a different time
Many systems offer to "remember a repeat" for you. Each one remembers for its own length of time:
| Who remembers a repeat | For how long |
|---|---|
| AWS Elemental MediaConvert client request token | 1 minute ("If you reuse a client request token within one minute of a successful request, the API returns the job details of the original request") |
Amazon SQS FIFO MessageDeduplicationId | 5 minutes |
DynamoDB TransactWriteItems ClientRequestToken | 10 minutes after the first request that uses it completes |
| Powertools for AWS Lambda idempotency, default | 1 hour |
| Stripe idempotency keys | may be removed once they are at least 24 hours old |
| The buyer, who closes the app and opens it tomorrow | days |
None of the first five is wrong. Each one is shorter than the last row.
What you'll be able to do after this page
- Tell a done, a not-done and an unknown outcome apart, and spot a success disguised as an error (Part 1).
- Design an idempotency key: who makes it, what it covers, and how it is tied to the request's meaning (Part 2).
- Claim a key, record the intent before the side effect, and answer every kind of repeat, including one from a crashed or paused owner (Part 3).
- Size a key's lifetime against the retry horizon, and stop an original that outlives its retry (Part 4).
- Make writes idempotent by design, and put the condition on the right record (Part 5).
- Explain what a queue's or a log's dedup covers and what the consumer's inbox must still catch (Part 6).
- Send a notification once per intent, and make an unavoidable repeat harmless (Part 7).
- Explain why a stream processor with exactly-once state still needs an idempotent sink, and what a restore must clean up (Part 8).
- Say what async multi-Region replication does to a claim, and fix it (Part 9).
- Compare idempotency keys with retry-free designs, locks, two-phase commit, sagas and broker exactly-once, on equal terms (Part 10).
- List every place along a request's path where a repeat is born, and the ID and guard for each (Part 11).
- Map all of it to AWS services, and name the look-alikes (Part 12).
You may have arrived from step 1.1 ("The Buyer Was Charged Twice") or 1.2 ("The PSP Call Timed Out. Did It Charge?") of the payments loop, step 1.2 of the mobile stock trading loop ("The Network Dropped After I Tapped Buy, So I Tapped Again"), step 2.4 of the ad-click aggregation loop ("A Worker Crash Replayed a Minute and Double-Billed") or step 1.3 of the notification loop ("A Retry Sent the Confirmation Twice"). This page is the "why" behind all of them.
Part 1. The core mental model
Before any mechanism, we need words for what a caller knows after it sends a request, and for what "once" can honestly mean. Everything later on the page follows from these.
Three outcomes of a request
A caller that sent a request ends up in one of three places:
| Outcome | What the caller saw | Examples | What to do |
|---|---|---|---|
| Done | A success reply | 200 OK, 201 Created | Nothing |
| Not done | A reply that says nothing happened | A validation 400, a card decline, a 409 that says "busy, try later" | Fix the request or try later |
| Unknown | No reply, or a reply that can't tell | A timeout; a connection reset after the request was sent; a 5xx from a system that may already have acted; an EIO from fsync | Record it as unknown. Resolve it by repeating with the same key, or by looking it up |
The third row is the dangerous one. A timeout says nothing about the server: the request may never have arrived, may have been done with the reply lost, or may still be running. Treat it as a failure and a buyer is charged with no order. Retry it blindly and the buyer may be charged twice.
The same holds inside a machine: the Write-Ahead Log loop primitive (Part 3) shows that an fsync error leaves the write's fate unknown, so the client's retry must be idempotent. This page is how.
Synthesizing vector architecture diagram...
What to notice: the caller sees the same thing, a timeout, in the second and third cases. Only the server knows which one happened, so the caller must ask in a way that is safe either way.
Success disguised as an error
Sometimes a retry does get a reply, and the reply is an error even though the first attempt worked:
| Retry | What comes back | What it really means |
|---|---|---|
A conditional PUT "create only if absent" (S3 If-None-Match: *, a DynamoDB put with attribute_not_exists) whose first reply was lost | 412 Precondition Failed or a failed condition | The object exists, probably because your first attempt created it. Read it and compare before calling this a failure (the key-value store loop, round 2) |
A Step Functions Standard StartExecution with the same name and input | While the execution runs: the original response. After it closes: ExecutionAlreadyExists | It started. AWS lets the name be reused 90 days after the execution closes. Express workflows aren't idempotent here |
The rule: an error on a retry is a question, not an answer. Read the resource before you decide.
Which HTTP methods are already safe to repeat
RFC 9110 (section 9.2.2) defines PUT, DELETE and the safe methods (GET, HEAD, OPTIONS, TRACE) as idempotent: sending one twice has the same intended effect as sending it once. POST is not, and "create a payment" is a POST. So a POST that creates something needs a key. Stripe, for example, accepts idempotency keys on every POST and says they have no effect on GET and DELETE, which are idempotent by definition.
Idempotent by definition is still about one request in isolation. Two different PUTs arriving in the wrong order is a different problem (Part 5).
Downstream windows
When we are the caller, the callee's own memory decides how we resolve an unknown. The PSP in our example keeps each key and its saved result for at least 24 hours, and compares the parameters of every repeat with the original:
- Inside the window: repeat the call with the same key and identical parameters. You get the original result back. Stripe saves the first result "regardless of whether it succeeds or fails", including
500errors. So a saved error keeps coming back for that key, and the way forward is to look the payment up. - After the window: the key may have been removed, and a repeat is then a new request. Never replay an old key; look the payment up by our own reference.
At-most-once, at-least-once, effectively-once
| Word | What it means | How you get it | What it costs |
|---|---|---|---|
| At-most-once | Never repeated; may be lost | Never retry | Lost payments and messages |
| At-least-once | Never lost; may be repeated | Retry until you hear back | Duplicates |
| Effectively-once | Every effect lands once, even though deliveries repeat | At-least-once, plus an effect that ignores repeats (a remembered key, or a write that is harmless to repeat) | A place to remember keys, and discipline at every hop |
"Exactly-once" is only true of an effect inside a store that records the key together with the effect, in one transaction. A delivery over a network can't be exactly-once: the sender can never be sure the receiver got it. So we say it hop by hop: "the charge is exactly-once at the PSP, the credit is exactly-once in the merchant database, the push is at-least-once with a replace ID".
Three rules the design never breaks
| Rule | What breaks without it |
|---|---|
| The client names the intent once, before the first send, and reuses the name on every retry | The server can't tell a retry from a new payment |
| The intent is recorded durably before the side effect | A crash after the PSP call leaves no trace, so nothing can resolve it |
| Every repeat, by anyone, carries the original ID | A relay, resolver or Region that mints a fresh ID defeats every dedup downstream |
Not this dedup
The word "dedup" also names mechanisms that have nothing to do with repeated requests. This page doesn't teach them:
- Content-addressed storage dedup: storing identical bytes once (the Drive loop's chunk store).
- URL seen-sets and Bloom filters: skipping pages a crawler has already fetched (the web crawler loop).
- Alert grouping: folding many alerts with the same labels into one page.
- Request collapsing: a cache or CDN turning many concurrent misses for the same object into one origin fetch. That merges reads that happen at the same moment, not a write repeated later.
The example we follow
Real payment systems handle millions of requests, which is too many to watch. So one story runs through the whole page: one payment, retried at every hop. Every trace on this page was produced by running a private reference simulator of this setup (the phone, a load balancer, two API servers, the key store, the ledger, the PSP, the resolver, the webhook, the relay, the queue, two workers, the push provider and a stream job), not worked out by hand.
| Setting | Our example | At real scale |
|---|---|---|
| Payment | Order 42, $30.00 (3,000 cents), buyer u1, merchant m1, one phone, one Region R1 (a standby Region R2 appears only in Part 9) | the payments loop: 50,000 a day in round 1, 20 million a day in round 2 |
| Client key | K1, a UUIDv7 (time-ordered, with 74 random bits) made by the app at the tap and written to the app's local journal before the first send | Stripe suggests V4 UUIDs or another high-entropy string, up to 255 characters, and no sensitive data in keys |
| Key scope | payments:u1:K1: the endpoint, the authenticated caller, the key | the payments loop scopes by endpoint (payments, refunds) |
| Fingerprint | H1: a hash of the request's meaning (order, amount, currency, payment method) | Stripe compares parameters |
| Durable guard | The ledger's payments row, inserted PENDING before any external call, with UNIQUE (caller_id, client_key) and the fingerprint; never deleted | the payments loop keeps UNIQUE (merchant_id, scope, idempotency_key) forever |
| Derived IDs | Payment ID pay_7, computed from the scope; PSP key pay_7:charge | the same rule everywhere |
| Event ID | pay_7:SUCCEEDED, derived from the payment and its new state, because this design has a standby Region (Part 9). We write it ev_7 for short | Change Streams & the Transactional Outbox (page 05) stores a random one: fine in one Region |
| Client timeout | 10 s from the first send; then retry with the same key and body | the mobile loops use 10 s (trading) and 30 s (news feed) |
App policy on 409 | Keep the key and the journal at SENDING, retry when the 10 s timer fires | polling with Retry-After is the alternative |
| Server deadline | 5 s from the load balancer's arrival stamp | the trading loop: 5 s against a 10 s app timeout |
| Key store | Table idem: state, fingerprint, owner, in_progress_until = arrival + 5 s, the answer, expires_at = completion + 3,600 s | Powertools' default 1 hour; the notification loop 24 hours; the payments loop 7 days |
| PSP | Keeps each key's result for 24 hours and compares parameters | Stripe |
| Resolver | Runs every second; replays payments that have been UNKNOWN for more than 5 s | the payments loop: every 30 s, after 10 s |
| Queue | FIFO, message group = merchant, dedup ID = the event ID, 5-minute dedup window, 30 s visibility timeout | SQS FIFO |
| Worker | One transaction: inbox insert, a credit entry for m1, a notification intent; then the push, with replace ID pay_7 | APNs apns-collapse-id, the Android notification tag |
| Stream job | Merchant daily totals from the ledger's change stream; checkpoint every 30 s; dedup horizon 2 hours | the ad-click loop |
m1 before the story | 12,000 cents |
The story in six beats: the reply is lost (Parts 2 to 4), the queue sees it twice (Part 6), the worker dies after acting (Parts 6 and 7), the stream job restores (Part 8), the buyer comes back two hours later (Part 4), and a Region fails in the middle of replication (Part 9). Each Part shows its own slice; the whole event table is in Part 13.
What to remember from Part 1
- A timeout is an unknown, not a failure: the work may have happened, and an "error" on a retry may hide a success.
- Retries give at-least-once; effects that ignore repeats turn that into effectively-once.
- Exactly-once holds only inside a store that records the key with the effect. Say it hop by hop.
Part 2. Idempotency keys
The server has to recognise a repeat. To do that, every copy of the request must carry the same name, and the name must mean "this intent", not "this attempt" and not "this order". This Part decides who makes the name, what it covers, and how it is tied to the request.
Who makes the key, and when
| Candidate key | What goes wrong |
|---|---|
The order ID (order 42) | The buyer's card is declined and they try another card: a genuinely new payment for the same order is refused as a "duplicate". And two tabs paying the same order at once look like one |
| A key the server generates and returns | The reply that carries it is exactly what gets lost. The client retries without it, and the server makes a second one |
| A hash of the request body | Two genuine $30 payments by the same buyer on the same day collapse into one |
| A random key the client makes once per intent, stores, then sends on every attempt | Works. A new intent (another card, a changed amount) gets a new key |
The key is born before the first send and written to the app's local journal, so an app that is killed mid-send still has it when it starts again.
The server could just generate the key and return it with the first response. Why doesn't that work?
What the key covers
The key is only unique within a scope: the endpoint, the authenticated caller, and the key. Ours is payments:u1:K1. Two reasons:
- Collisions across callers can't hurt. Two buyers who happen to send the same key never see each other's answers.
- Keys can't be used to probe. A caller can only ever replay answers to its own requests.
Here the caller is the buyer. In a server-to-server API, such as a merchant calling a payments platform, the caller is the merchant.
Binding the key to the request's meaning
A client bug can reuse a key for a different request: the same K1 with $35.00 instead of $30.00. Replaying the $30 answer would be wrong, and so would running the $35 payment. So the server stores a fingerprint with the key: a hash of the request's meaning, not of its bytes.
textSCOPE(request) = endpoint + ":" + authenticated caller + ":" + Idempotency-Key header FINGERPRINT(request) = hash of the canonical meaning fields, in a fixed order: order id, amount in minor units, currency, payment method -- not the raw body: a reordered JSON body, a new client -- timestamp field or an extra header is the same request on a repeat with the same scope: stored fingerprint == new fingerprint -> the same request: answer it (Part 3) stored fingerprint != new fingerprint -> 422 "key reused with a different request", nothing runs
Side trace 3d, from the simulator: A2 is sent at 10.00 with K1 but an amount of 3,500. Its fingerprint H2 differs from the stored H1, so the server answers 422 and nothing changes: no claim, no PSP call. (Some APIs use 409 here; Stripe errors on mismatched parameters.) The same run confirmed the other direction: a body with the same fields in a different order, plus an extra client timestamp field, gives the same fingerprint H1.
Keys must not contain personal data (Stripe's advice: no email addresses or personal identifiers). They end up in logs, tables and support tools.
Derive what others must re-create; store what you send
The key protects one hop. Every later hop needs its own ID, and the question is how each one is made:
- Derive an ID if someone else may have to re-create it. A resolver replaying a charge, another Region taking over, or a re-run parser must produce the same ID as the first attempt, or the next hop sees a stranger. So
pay_7is computed from the scope, and the PSP key ispay_7:charge. If a second attempt is legitimate (a failover to a second PSP after a definite failure), it gets a new suffix, one per attempt per PSP:pay_7:attempt-2:psp_b. The cost: a derived ID leaks structure, so hash it if it must be opaque. - Store an ID if it is written down before its first send. An outbox row's event ID (Change Streams & the Transactional Outbox, page 05) is stored in the same transaction as the payment, and every re-send reads it back from the row. So a random ID would do. Ours is derived anyway, as
pay_7:SUCCEEDED, because in Part 9 another Region has to re-create the event. - Use natural keys where they exist: a
click_id, atrade_id, a chat app'sclient_msg_id.
Synthesizing vector architecture diagram...
What to notice: the notification intent's key (ev_7 + channel) and the replace ID (pay_7) are two different IDs with two jobs. The intent key stops us from sending twice; the replace ID tells the phone that a second copy is the same notification.
Standards: there is no standard idempotency header yet. The IETF draft "The Idempotency-Key HTTP Header Field" reached revision 07 and is now an expired draft. The common practice this page describes (a client-made key in an Idempotency-Key header, scoped to the caller, bound to the request) is what payment APIs do in practice.
Trace: the tap and the claim
| # | t (s) | Event |
|---|---|---|
| 1 | 0.00 | Tap. The app creates K1, writes {K1: SENDING} to its journal, and sends A1 over Wi-Fi: POST /payments, Idempotency-Key: K1, body {order 42, 3000, USD, pm_1} |
| 2 | 0.02 | A1 reaches the load balancer (arrival stamp 0.02) and api-1. Claim idem[payments:u1:K1] only if absent: IN_PROGRESS, fingerprint H1, owner A1, in_progress_until = 0.02 + 5 = 5.02. Derive pay_7. Insert pay_7 as PENDING with (u1, K1) and H1, before any external call |
The request on the wire:
httpPOST /payments HTTP/1.1 Host: api.example.com Authorization: Bearer <u1's token> Idempotency-Key: 0192f0c4-5b7e-7c1a-9d3e-2f4b8a6c1e07 Content-Type: application/json {"order_id": "42", "amount": 3000, "currency": "USD", "payment_method": "pm_1"}
What to remember from Part 2
- The client names the intent once, before the first send, and reuses the name on every retry.
- Scope the key to the caller, and bind it to a hash of the request's meaning: the same key with a different request is a bug, answered with
422. - Derive every ID someone else may have to re-create; an ID you store before sending may be random.
Part 3. The key record and the durable guard: claim, record, run, answer
Two copies of the same request can run at the same time. A finished request must answer its repeats the same way. And a crash right after the PSP call must leave something behind that says "a charge may be out there". The naive answers fail: "check whether a payment exists, then create one" races (two copies both see nothing), and "write the payment row once the PSP answers" leaves no trace if we crash in between. The fix is the rule the Write-Ahead Log loop primitive is built on: log first, then act. Here, record the intent before the side effect.
Record the intent before the side effect
textHANDLE(request) -- on the API server 1. if now > arrival stamp + deadline: reject; run nothing (Part 4) 2. claim idem[scope] only if it is absent (or its expires_at has passed): state IN_PROGRESS, fingerprint, owner = this request, in_progress_until = arrival stamp + deadline claim failed -> ANSWER A REPEAT (below) 3. payment_id = derive(scope) insert the payments row: payment_id, caller, key, fingerprint, state PENDING -- UNIQUE (caller_id, client_key); this happens BEFORE any external call unique conflict -> read that row and answer from it (a VOID row -> 409; another fingerprint -> 422) 4. call the PSP with key payment_id + ":charge" 5. the reply arrives in time: one ledger transaction: PENDING -> SUCCEEDED (compare-and-set), journal entries, outbox row complete idem, only if owner = me: store the final answer the deadline comes first: PENDING -> UNKNOWN (compare-and-set) complete idem, only if owner = me: store a POINTER to payment_id 6. completion failed (the owner changed) -> read idem and answer what it holds
The claim, as a DynamoDB conditional put (times in the story's seconds; real items use epoch time):
json{ "TableName": "idem", "Item": { "pk": { "S": "payments:u1:K1" }, "state": { "S": "IN_PROGRESS" }, "fp": { "S": "H1" }, "owner": { "S": "A1" }, "in_progress_until": { "N": "5.02" } }, "ConditionExpression": "attribute_not_exists(pk) OR expires_at < :now", "ExpressionAttributeValues": { ":now": { "N": "0.02" } } }
The durable guard, in the ledger:
sqlCREATE TABLE payments ( payment_id TEXT PRIMARY KEY, -- derived from the scope: pay_7 caller_id TEXT NOT NULL, -- u1 client_key TEXT NOT NULL, -- K1 fingerprint TEXT, -- H1; NULL only on a VOID row amount BIGINT, -- 3000; NULL only on a VOID row currency CHAR(3), state TEXT NOT NULL, -- PENDING, UNKNOWN, SUCCEEDED, FAILED, VOID psp_charge TEXT, -- ch_1 once known UNIQUE (caller_id, client_key) -- the durable guard: lives as long as the row );
The states of a key record
Synthesizing vector architecture diagram...
What to notice: "absent" includes "expired": a reader treats a record whose expires_at has passed as absent, whether or not it has been deleted yet. When the key record expires, the ledger row (including a VOID row) is still there, and it is what answers from then on (Part 4).
What gets stored: final answers and pointers
| Answer | Example | Stored as | Replayed as |
|---|---|---|---|
| Final success | 200 {pay_7, SUCCEEDED} | The whole answer | Byte for byte, with a replay marker header |
| Final rejection | A decline, a 422 | The whole answer, rejections included | Byte for byte: "same request, same result" |
| Not final yet | 202 PROCESSING | A pointer to pay_7 | Rendered from pay_7's current state, so a replay after the payment settles says 200 |
A pointer matters because the state behind it changes. If we stored 202 PROCESSING byte for byte, a retry an hour later would still hear "processing" about a payment that finished long ago. When the payment does become final, the resolver also rewrites the stored answer to the final one.
Trace: a concurrent duplicate, a lost reply and a retry
| # | t (s) | Event |
|---|---|---|
| 3 | 0.05 | api-1 sends the charge to the PSP with key pay_7:charge |
| 4 | 0.30 | The PSP charges the card (ch_1) and saves the result under pay_7:charge. Its reply is delayed: a gray failure |
| 5 | 0.45 | The phone switches from Wi-Fi to cellular, and A1's Wi-Fi connection dies. The phone's HTTP library resends A1′ (same key, same body) on a new cellular connection. It reaches api-2: the claim fails, the record is IN_PROGRESS, the fingerprint matches → 409, Retry-After: 1. The app keeps waiting for its 10 s timer |
| 6 | 5.02 | api-1 reaches its deadline (0.02 + 5) with no reply. The call was sent, so the outcome is unknown, not abandoned: pay_7 PENDING → UNKNOWN. Then idem → COMPLETED with a pointer to pay_7, expires_at = 5.02 + 3,600 = 3,605.02. api-1 writes 202 PROCESSING |
| 7 | 5.03 | The 202 has nowhere to go. It belongs to A1's Wi-Fi connection, which died at 0.45. A reply can only travel back on the connection its request came in on |
| 8 | 5.60 | The PSP's reply arrives on api-1's connection to the PSP, closed at the deadline: lost. Only the resolver or the PSP's webhook can tell us now |
| 9 | 10.00 | The app's 10 s timer fires. Retry A2 (same key, same body) over cellular → api-2: COMPLETED, fingerprint matches, the answer is a pointer → rendered from pay_7's current state (UNKNOWN) → 202 PROCESSING with a replay marker. The journal says UNKNOWN; the app shows "Confirming your payment" and polls |
Snapshot P1, right after event 5:
Synthesizing vector architecture diagram...
After event 5: the card has been charged, the ledger knows a charge may be out there (PENDING), and the phone knows nothing.
Snapshot P2, after event 9 (t = 10.00):
| Lane | State |
|---|---|
| Phone | Journal K1: UNKNOWN; shows "Confirming your payment" |
| Key store | COMPLETED, fingerprint H1, answer = pointer to pay_7, expires_at 3,605.02 |
| Ledger | pay_7 UNKNOWN, no charge ID recorded, no outbox row |
| Downstream | PSP holds ch_1 under pay_7:charge; queue empty; m1 = 12,000 |
A1 is still running when A1′ arrives at 0.45 with the same key and body. Should api-2 wait, run it again, or answer something?
Every answer to a repeat
| What the claim finds | Answer |
|---|---|
| Nothing (or an expired record) | Run it: this is the first copy the server has seen, or the window has passed and the durable guard will decide (Part 4) |
COMPLETED, same fingerprint | Replay: a final answer byte for byte, or a pointer rendered from the payment's current state |
Any state, different fingerprint (a VOID row has no fingerprint and answers 409) | 422: the key was reused for a different request; nothing runs |
IN_PROGRESS, not yet expired | 409 + Retry-After (two stores); in one PostgreSQL transaction, the insert waits for the first transaction, then conflicts |
IN_PROGRESS, past in_progress_until + margin | Take over (next section) |
VOID | 409: the key was voided by a resolve call; nothing runs (Part 4) |
A crashed or paused owner
A request that claimed a key can crash, or freeze (a long garbage-collection pause, a paused VM), and leave the record IN_PROGRESS. A later repeat may take over, with three safeguards:
- One clock decides expiry. Either a single SQL database's
now()for rows in that database (the payments loop does this), or the claimant's own clock plus a margin larger than the worst clock offset between machines. We use 1 s. Leases, Fencing Tokens & Distributed Locks (Part 5) explains why the margin must exceed the offset. - The takeover is conditional on the owner still being the one the claimant read, and it checks the durable record first, because the old owner may have got further than the key record says.
- Completion is conditional on
owner = me. A paused owner that wakes up can't overwrite the new owner's answer. It reads the record instead.
From the simulator, three variants of the main line:
| Side trace | What happened to api-1 | What A2 finds at 10.00 | What A2 does | Result |
|---|---|---|---|---|
| 3a | Committed pay_7 UNKNOWN at 5.02, then crashed before completing idem | IN_PROGRESS, in_progress_until 5.02; 10.00 > 5.02 + 1 → takeover (owner A1 → A2) succeeds | Reads pay_7: UNKNOWN → completes idem with a pointer, answers 202 | The resolver settles it at 11.00, as in the main line. One charge |
| 3b | Crashed at 1.00, after its PSP call went out at 0.05 | Same takeover | Reads pay_7: PENDING, so a charge may be out: resolves with pay_7:charge and identical parameters, gets the saved ch_1, moves PENDING → SUCCEEDED, answers 200 | One charge |
| 3c | Paused, not dead, from 5.02 (after UNKNOWN) to 12.00 | Same takeover | As in 3a. At 11.00 the resolver makes the answer final. At 12.00 api-1 wakes; its completion "only if owner = A1" fails (owner is A2), so it reads the record: COMPLETED, 200 {pay_7, SUCCEEDED} | One charge |
In 3c, a paused owner is a zombie: it still believes it owns the request. Its PSP call, had it made one after waking, would have been harmless only because the PSP key is derived (pay_7:charge), so the PSP replays instead of charging. Stopping a zombie's writes to your own stores is the fencing-token problem; the Leases page owns it.
In our simulator the resolver only picks up UNKNOWN payments, so in 3b the PENDING row waits for the phone's retry. A production resolver also sweeps PENDING rows whose deadline plus margin has passed, so a crashed owner's payment is settled even if the phone never comes back; its compare-and-set on the state makes it safe to race with a takeover.
One store or two, and what if the key store is down
| One store: the key row in the same database transaction as the payment (the wallet, hotel and trading loops) | Two stores: a fast key store in front of the durable guard (the payments loop's round 2 on DynamoDB and Aurora, the notification loop) | |
|---|---|---|
| Consistency | The key and the payment commit together; they can't disagree | They can disagree: an expired or lost key item while the row exists, or an IN_PROGRESS item with no row after a crash. The durable guard wins |
| Hot-path cost | One short transaction before the PSP call (the payments loop: about a millisecond) | Two conditional writes before the PSP call |
| Concurrent repeat, default answer | The second insert waits for the first transaction, then conflicts: "wait, then answer" | 409 + Retry-After |
| Scaling | Every claim, including abusive retry loops, lands on the ledger's one writer | Claims go to a store that scales out; the ledger sees only real work |
| When the key store is down | Not applicable | Fail closed for money (refuse, don't guess), or fall back to claiming directly against the durable guard's unique key. The guard alone is enough for correctness; you lose only stored-answer replay and the fast 409 |
Finish what the first attempt started
A repeat that is recognised as a duplicate must not stop at "already done". The first attempt may have died between two steps, so the duplicate finishes every step the first attempt may not have finished. The chat loop's retry re-sends the acknowledgement and re-publishes the message; in our story, A3 at 7,200 rewrites the expired key record, and in Part 7 the second worker sends the notification the first may never have sent.
What to remember from Part 3
- Write the intent durably before the first external call; every repeat then runs into it.
- A repeat while in progress waits or gets
409; after completion it gets the stored final answer, or the payment's current state if the answer wasn't final. - A takeover uses one clock and checks the durable record; a zombie's completion fails on the owner check.
Part 4. Time: how long to remember, and the late original
The key record is a window: in our example it forgets K1 an hour after completing. Phones retry for longer than that. And there is a second, sneakier timing problem: an original request that runs after its client has given up on it.
When a late original can hurt
A late original can hurt only when the client stops using a key and makes a new one, as when the buyer taps Pay again. Retries with the same key are already safe: whichever copy comes first claims the key, and the rest are answered. The danger is a first request still sitting in a queue somewhere when the second key starts its own payment. Two tools close that gap: a deadline on the server, and a "did it happen?" call that voids the old key.
The retry horizon
The retry horizon is the longest time after the first send at which a retry can still arrive:
| Client | How long it retries |
|---|---|
| An AWS SDK call in standard retry mode | 3 attempts in total by default (1 request + 2 retries), within seconds; under AWS's updated retry behaviour DynamoDB clients use 4, and older SDK defaults vary |
| A browser checkout page | Until the tab is closed: minutes |
| A mobile app with a local journal | Days to weeks: the mobile news feed loop (step 2.2) holds actions while a phone stays offline |
| Stripe delivering webhooks to us | Up to three days in live mode |
| A server-to-server caller with a durable job queue | As long as its own retry budget: hours to days |
Formula 1:
Our key record remembers K1 until 3,605.02. The buyer comes back at 7,200: a horizon of at least 7,200 s against a window of 3,600 s, so the window fails. The durable guard remembers K1 as long as the payment row exists, so the guard passes. Every window on the Part 0 table (1 minute to 24 hours) is shorter than a phone's horizon. That is fine, as long as each window is treated as an accelerator and correctness sits on the durable guard.
Expired means expires_at < now, checked in the read path. DynamoDB's TTL deletes expired items "within a few days of their expiration time", so an expired item can still be there. TTL is cleanup, never the window.
Trace: the buyer comes back
| # | t (s) | Step |
|---|---|---|
| 20a | 7,200.00 | The buyer opens the app. The journal says K1: UNKNOWN, so the app resends A3 with the same key and body |
| 20b | 7,200.00 | idem[payments:u1:K1] has expires_at 3,605.02 < 7,200, so it counts as absent even if TTL hasn't deleted it. The claim succeeds |
| 20c | 7,200.00 | pay_7 is derived again. Its insert hits UNIQUE (u1, K1) before any PSP call. The row says SUCCEEDED (ch_1) and its fingerprint H1 matches |
| 20d | 7,200.00 | idem is rewritten COMPLETED (new expires_at 7,200 + 3,600 = 10,800); the answer is 200 {pay_7, SUCCEEDED} with a replay marker. A different fingerprint would have got 422 |
Snapshot P6:
Synthesizing vector architecture diagram...
What to notice: the key store forgot K1 an hour ago and it didn't matter: the ledger row answered, and nothing reached the PSP.
Branch X, the same event without the durable guard. Take away the unique key and mint the payment ID at random, and the simulator gives: the expired key looks new, the server creates pay_9, and its PSP key pay_9:charge is new to the PSP, which charges ch_2 at 7,200.28. The buyer has paid $60.00 for one order, and was told "succeeded" for both.
The buyer comes back after two hours; the key store forgot K1 an hour ago. What stops a second charge?
The deadline
In the opening story, A1 reached the load balancer at 0.02 and waited in api-1's queue until 12.00, while the phone gave up at 10.00 and the buyer tapped again with a new key. The in-progress lease can't help: A1 hadn't claimed anything yet. What stops it is a deadline measured from when the request arrived:
- The load balancer stamps the arrival time on the request.
- The server refuses to start a request whose stamp is more than 5 s old, and when a running request reaches 5 s it starts no new side effect and commits nothing new. A call already sent becomes
UNKNOWN, never "abandoned". - The server checks the stamp on its own clock, so it needs a margin for the offset between the load balancer's clock and its own (the Leases page, Part 5).
Formula 2:
With the stamp at the load balancer, the only wait before it is the network trip: 5 + 0.02 = 5.02 s < 10 s. By the time the phone's timer fires at 10 s, any copy of A1 has either claimed K1 already (so a "did it happen?" call sees it) or will never run. If the stamp is taken only when a server thread picks the request up, the wait before the stamp includes the load balancer's, the connection and the thread-pool queues. In our story that wait was 12.00 − 0.02 = 11.98 s, and 5 + 11.98 = 16.98 s > 10 s. Then you must bound those queues too: shedding requests that waited more than 2 s gives 2 + 5 = 7 s < 10 s.
On AWS: an Application Load Balancer adds an X-Amzn-Trace-Id header whose ID fields carry the epoch time in seconds (8 hex digits). If the client already sent a Root field, the load balancer inserts a Self field and keeps the client's Root; read the load balancer's own field, not one a client could send. The time is in whole seconds, so it can read up to 1 s early, which only makes the deadline stricter. Also note that the load balancer won't end a slow request when the phone gives up: its idle timeout defaults to 60 s. The server's own deadline is what stops the work.
Synthesizing vector architecture diagram...
What to notice: in the second half, two separate rules each stop A1. The deadline rejects it before it claims anything; and even without the deadline, A1's claim would find K1 as VOID and get 409. The simulator ran both: with only the void, A1 gets 409 at 12.00; and after the key record's hour has passed, A1's insert hits the VOID row under UNIQUE (u1, K1) and still gets 409, with no PSP call.
Resolve, don't just look
The app's "did it go through?" call must change the state, not just read it:
textRESOLVE(caller, key) -- the app's "what happened to my request?" 1. read idem[scope] and the payments row for (caller, key) 2. found: return its state ("still confirming" for PENDING or UNKNOWN; else the final answer) 3. absent: insert-if-absent in BOTH stores: idem[scope] = VOID a payments row (caller, key, state VOID) under UNIQUE (caller_id, client_key) return "not found, and it never will be": the app may now make a new key 4. a late original that arrives later finds VOID and gets 409; nothing runs
The insert-if-absent races fairly with a late original's claim: whichever gets the unique row first wins. If the original wins, the resolve call reports "in progress", and the app keeps waiting instead of making a new key.
One more hole: an app reinstall loses the journal and its keys. The trading and news-feed loops keep, per device, the last accepted action ID, so a reinstalled app can ask "what was my last order?" before it lets the user act again.
When a compensation times out
The same rules apply to the steps that undo a payment. Suppose an order fails after the charge and we refund. The refund call to the PSP times out. That is an unknown, not a failure: record it as such, and retry it with its own derived key, pay_7:cancel:refund, so every retry is the same refund. Keep retrying with backoff; after the retry budget, also alert a human, but keep the state UNKNOWN rather than guessing. The nightly reconciliation against the PSP's records has the last word. Nothing is lost permanently because the intent to refund was recorded before the first call and every attempt carries the same key. This answers the first question of the drill The Flight Booking That Charged Without a Seat; Part 10 answers the second.
What to remember from Part 4
- A key store's window is an accelerator; the durable record decides.
- The server's deadline, counted from arrival, must end before the client can give up and make a new key.
- A "did it happen?" call must void the key, not just look.
Part 5. Idempotent by design: absolute writes and conditions on the durable record
A key table in front of every write is heavy. Many writes can be made harmless to repeat by their shape alone, with no key table at all. And where a condition is needed, it matters which record carries it.
Write the result, not the change
| Operation | Idempotent form | What the idempotent form can't do |
|---|---|---|
Add to a total: ADD 3000 | Absolute write: SET total = 15000, computed from state that already counted the event once | Needs the full value; two writers computing it separately must be ordered by a version |
| Count something | Set-if-larger, or one item per window (clicks#m1#10:05) written with its absolute count | Set-if-larger can never lower a value, so corrections need a version |
| Add to a set | Set semantics: ZADD of an exact member, or insert-if-absent of a natural key (click_id, trade_id) | "Only if greater" writes such as ZADD GT always add a missing member, so they don't stop a stale event from re-adding one |
| Change a state | Compare-and-set: UNKNOWN → SUCCEEDED only if the state is still UNKNOWN | Needs a state machine with the allowed moves written down |
| Apply a change from a stream | Versioned upsert: apply only if version > stored version | Versions must come from one counter per key |
| Write output files | Fixed output keys, each written as a whole object | Nothing, as long as each object is written whole |
ev_7 credits m1 with ADD 3000. It's delivered twice. Name two ways to make the second delivery harmless.
State transitions, and the late webhook
Our payment moves only along allowed transitions (PENDING → UNKNOWN → SUCCEEDED, and so on), and every move is a compare-and-set. That makes a late or repeated answer harmless:
| # | t (s) | Event |
|---|---|---|
| 11 | 11.00 | The resolver runs (every second; at 10.00 pay_7 had been UNKNOWN for 10.00 − 5.02 = 4.98 s, at 11.00 for 5.98 s > 5 s). It replays the charge with pay_7:charge and identical parameters; the PSP returns the saved ch_1. One ledger transaction: pay_7 UNKNOWN → SUCCEEDED (compare-and-set), journal entries, outbox row ev_7. Then idem's answer is rewritten to the final 200 {pay_7, SUCCEEDED} |
| 13 | 11.50 | The PSP's webhook payment_intent.succeeded for ch_1 arrives. Its event ID evt_1 is inserted into webhook_events (new). The compare-and-set UNKNOWN → SUCCEEDED finds SUCCEEDED: nothing changes |
The order doesn't matter. The simulator also ran the webhook first, at 10.50: it made the same move (UNKNOWN → SUCCEEDED, outbox row ev_7), a second delivery of evt_1 at 10.90 was dropped as a duplicate, and at 11.00 the resolver found nothing UNKNOWN to do. The stored answer was still the pointer, so it rendered SUCCEEDED without being rewritten.
Versions from one counter
Idempotence alone doesn't cover order: an old message arriving after a newer one. Each copy is harmless on its own, but applying the old one last undoes the new one. The fix is a version from one counter, compared on write.
Branch V is the mobile news feed loop's case (step 2.2). Device D1 keeps a counter in its journal. It sends LIKE (counter 41), then UNLIKE (42); the LIKE is resent and arrives last. From the simulator:
| Arrives | By device counter: apply if stored < incoming | By phone timestamp (the phone's clock was stepped back 3 s between the taps) |
|---|---|---|
LIKE 41 (phone time 10:00:05) | 0 < 41 → applied; liked | applied; liked |
UNLIKE 42 (phone time 10:00:03) | 41 < 42 → applied; unliked | 10:00:03 is older than 10:00:05 → ignored; still liked |
LIKE 41 again | 42 ≥ 41 → DUPLICATE; stays unliked | ignored; still liked |
| Final | Unliked: correct | Liked: the user's unlike is lost |
Never order by a timestamp made on the phone. Phone clocks are set by users, carriers and time daemons, and they jump. The counter must also outlive anything it's stored next to: keep the per-device high-water mark on an item that never expires.
Condition on the durable record
A condition protects you only if it sits on a record that exists for as long as the effect does. Put it on an item that can move or expire, and a retry that comes after the move re-creates the effect.
Branch M is the email loop's case (round 2, delivery): a parser, fed by SQS, stores message m9, and SQS redelivers the parse after the user has moved m9 from Inbox to Archive. From the simulator:
Synthesizing vector architecture diagram...
What to notice: both writes are "put if absent". The left one conditions on a folder item that the move deleted, so the retry passes and m9 reappears in Inbox. The right one conditions on the locator M#m9, which never moves and becomes a tombstone when m9 is deleted, so the retry fails even after a delete (the simulator's second retry, after deletion, also did nothing).
The same rule appears across the loops: the payments loop's retry conditions on the payment row, not on the key store's item that expires (Part 4); YouTube's "ready" handler conditions on the video's state, not on a job token that MediaConvert forgets after a minute. When a condition fails, read what's there and answer from it; a failed condition on a retry is usually a success disguised as an error (Part 1).
What to remember from Part 5
- Write the result, not the change: absolute values, set-if-larger, insert-if-absent.
- Put the condition on the record that never moves or expires.
- Versions from one counter stop an old repeat from undoing a newer change; phone clocks can't.
Part 6. Queues and logs: dedup windows, redelivery and the inbox
Between the ledger and the worker, repeats are born in two places: the producer sends the same event twice (a relay that crashed after publishing), and the broker delivers the same message twice (a worker that died before deleting it). Brokers help with the first, only for a while. The second is always the consumer's job.
Two places repeats are born
Synthesizing vector architecture diagram...
What to notice: the producer-side window catches only fast repeats; everything it misses, and every redelivery, ends at the consumer's guard. The guard is the one that can't be skipped.
The FIFO dedup window, and what keeps the ID stable
Facts from the SQS documentation:
- A FIFO queue drops a message whose
MessageDeduplicationIdit has seen in the last 5 minutes: the send is "accepted successfully" but the message isn't delivered again. - SQS keeps tracking the ID even after the first copy was received and deleted.
- Content-based dedup uses a SHA-256 hash of the message body (not its attributes) as the ID. So a body that carries an attempt count or a send timestamp makes every repeat look new.
- With high-throughput mode, the dedup scope is the message group, so a re-send must keep the same group ID as well as the same dedup ID.
- SNS FIFO topics have the same five-minute window for publishes.
The dedup ID is only as stable as the process that picks it. Ours is the event ID stored in the outbox row, read back by every re-send. The outbox relay (Change Streams & the Transactional Outbox, Part 3: claim, publish, delete) re-sends an undeleted row with the ID it finds there.
Trace: the relay crashes twice
| # | t (s) | Event |
|---|---|---|
| 12 | 11.20 | The relay claims the ev_7 row and sends it, dedup ID ev_7: accepted, window until 11.20 + 300 = 311.20. The relay crashes before deleting the row |
| 14 | 14.00 | The relay restarts, claims the still-undeleted row and sends it again with the same dedup ID, 14.00 − 11.20 = 2.80 s into the window: the queue acknowledges but doesn't enqueue. The relay crashes again before deleting the row |
Snapshot P3, after event 14:
Synthesizing vector architecture diagram...
After event 14: two sends, one message. The undeleted outbox row will be sent again, and nothing upstream can stop that.
The idempotent producer, and why a restart defeats it
A log such as Kafka offers an idempotent producer (enable.idempotence, on by default in current clients). The broker gives the producer an ID, and the producer numbers its batches per partition; a batch with a sequence number the broker has already written is dropped. Kafka's own description: "the producer will ensure that exactly one copy of each message is written in the stream". The catch is the scope: one producer session.
Branch K, from the simulator. The relay publishes ev_7 to a log instead of a queue:
| Step | What happens | Log afterwards |
|---|---|---|
| 1 | Session 1 gets producer ID 17 and sends ev_7 as batch sequence 0. The broker writes it; the acknowledgement is lost | 17/0 ev_7 |
| 2 | The library retries sequence 0: the broker drops it as a duplicate | 17/0 ev_7 |
| 3 | The relay crashes before deleting the outbox row and restarts without a transactional ID: it gets producer ID 18 | |
| 4 | It re-sends the row as 18/0: accepted, because producer 18 has sent nothing before. A new session is a stranger | 17/0 ev_7, 18/0 ev_7 |
Transactions in the log
A transactional ID makes the producer's identity survive restarts: Kafka documents that it "enables reliability semantics which span multiple producer sessions", and a new session with the same ID first completes or aborts whatever the previous session left open. Output records and consumed offsets can then commit together. Readers must set isolation.level=read_committed; the default, read_uncommitted, also returns records from aborted transactions.
Branch K again, with transactional ID relay-1:
- Session 1 commits
ev_7, then crashes before deleting the outbox row. Session 2 starts (epoch 1), re-reads the row and commits it again. Aread_committedreader seesev_7twice. The transaction covered the log; the outbox row that caused the re-send is outside it. - Session 1 crashes before committing. Session 2's start aborts session 1's open record;
read_committedseesev_7once, whileread_uncommittedsees it twice.
So a transactional log removes some repeats, but an outbox re-send is a new record. Only the consumer's guard catches it.
The consumer's two guards
A consumer that applies an effect must recognise a repeat itself. Two ways:
- An inbox keyed by event ID, written in the same transaction as the effect. This page's choice, for effects like credits that aren't naturally ordered.
- Apply by version, for an ordered stream of changes to one key (a cache or an index following a row). Change Streams & the Transactional Outbox (Part 3) covers that.
The worker's transaction:
sqlBEGIN; INSERT INTO processed (event_id) VALUES ('pay_7:SUCCEEDED') ON CONFLICT (event_id) DO NOTHING; -- 1 row: first time; 0 rows: a repeat -- only if 1 row was inserted: INSERT INTO credits (merchant_id, cents, event_id) VALUES ('m1', 3000, 'pay_7:SUCCEEDED'); INSERT INTO notif (event_id, channel, state) VALUES ('pay_7:SUCCEEDED', 'push', 'PENDING') ON CONFLICT (event_id, channel) DO NOTHING; COMMIT; -- then, whether or not a row was inserted: if notif is PENDING, send the push (Part 7)
ON CONFLICT DO NOTHING matters in PostgreSQL: a plain INSERT that hits the unique key raises an error, and the error aborts the whole transaction. With DO NOTHING, the insert reports 0 rows and the worker decides what to skip.
Trace: the worker dies after acting
| # | t (s) | Event |
|---|---|---|
| 15 | 14.50 | Worker w1 receives ev_7; it is hidden from other workers until 14.50 + 30 = 44.50. One transaction: processed[ev_7] inserted (1 row), credit m1 +3,000 (12,000 → 15,000), notif[ev_7, push] = PENDING |
| 16 | 14.60 | w1 sends the push "Payment of $30.00 confirmed" with replace ID pay_7 |
| 17 | 14.65 | w1 crashes before marking the notification sent and before deleting the message |
| 18 | 44.50 | The visibility timeout ends and ev_7 goes to w2. The processed insert returns 0 rows → skip the credit. notif is still PENDING → at 44.60, send again with the same replace ID. Mark SENT, delete the message |
Synthesizing vector architecture diagram...
What to notice: the inbox stopped the second credit, and the notification record made w2 finish the one step w1 may not have finished. Whether the first push ever reached the phone, w2 can't know, so it sends again with the same replace ID.
Snapshot P4, after event 18:
Synthesizing vector architecture diagram...
After event 18: one credit, one notification on screen, and one outbox row still waiting to be sent a third time.
A re-send carries the original ID
| # | t (s) | Event |
|---|---|---|
| 19 | 431.20 | The relay comes back (or another relay instance claims the row) and re-sends the still-undeleted ev_7 row with its stored ID. The window ended at 311.20, 431.20 − 311.20 = 120 s earlier, so ev_7 is delivered. w1 (restarted)'s inbox insert returns 0 rows → no credit; notif is SENT → no push. The relay deletes the row |
| Branch Q | 431.20 | The same event without the inbox: w1 credits m1 again, 15,000 → 18,000: +6,000 for one $30.00 payment. (Without an inbox at all, event 18 would also have credited: the simulator ends at 21,000) |
The relay was down for 7 minutes and then re-sent ev_7 with the same dedup ID. Did the queue drop it?
Webhooks from outside are the same case. The PSP delivers each webhook at least once (Stripe retries for up to three days), so the webhook handler keeps an inbox keyed by the sender's event ID (evt_1 in event 13) and applies each event as an allowed state move.
Shared limits on this path. The FIFO message group is the merchant, so one busy merchant's events are serialised, and a message in flight (up to its 30 s visibility timeout) holds back the rest of its group. The inbox shares the consumer database's write path with the credits, and if the consumer also keeps a running balance per merchant, every credit writes that one item: a DynamoDB partition gives at most 1,000 write units a second, and a SQL row takes one lock at a time. Plan for the biggest merchant, not the average one.
What to remember from Part 6
- A queue's or producer's dedup window stops fast repeats only; the consumer's guard stops the rest.
- The idempotent producer covers its own retries within one session; a restarted producer is a new sender.
- Whoever re-sends must re-send the original ID; a new ID defeats every dedup downstream.
Part 7. Side effects you can't take back: notifications
A push, an email or an SMS leaves our systems the moment we send it. It can't join our transaction, and the providers take no idempotency key. So "send exactly once" is impossible: the honest goals are send once per intent, and make an unavoidable repeat replace instead of stack.
Why providers can't help
| Provider call | Takes an idempotency key? | Can a repeat replace the first? |
|---|---|---|
Amazon SES v2 SendEmail | No | No: a second email is a second email |
AWS End User Messaging SMS v2 SendTextMessage | No | No |
| Apple Push Notification service (APNs) | No | Yes: the apns-collapse-id header (at most 64 bytes). Apple: "When sending the same notification more than once, use the same value in this header to merge the requests" |
| Firebase Cloud Messaging (FCM), Android | No | Yes, with the Android notification tag: "If specified and a notification with the same tag is already being shown, the new notification replaces the existing one in the notification drawer" |
FCM's collapse_key looks like the answer and isn't. It lets FCM keep only the latest message of a group while a device is offline, "so that only the last message gets sent when delivery can be resumed", and FCM keeps at most four different collapse keys per device. It never touches a notification already on the screen.
Intent, send, mark
textIN THE EFFECT'S TRANSACTION (the worker, Part 6) insert-if-absent notif[event_id, channel] = PENDING AFTER THE COMMIT, and again on every redelivery of the event: if notif[event_id, channel] is SENT: skip if PENDING: send with replace id = the thing announced (pay_7) then set notif[event_id, channel] = SENT
The intent row makes "should I send?" a question with a durable answer. A crash between the send and SENT is the one case that sends twice, and the replace ID makes that second copy land on top of the first.
| # | t (s) | Event |
|---|---|---|
| 15 | 14.50 | notif[ev_7, push] = PENDING, in the credit's transaction |
| 16 | 14.60 | Push sent, replace ID pay_7: the phone shows it |
| 17 | 14.65 | w1 crashes before writing SENT |
| 18 | 44.60 | w2 finds PENDING → sends again with replace ID pay_7: the phone replaces the notification (it may alert again; if the buyer already dismissed it, it appears again). Then SENT |
w1 sent the push and died before writing SENT. What should w2 do, and what makes that safe?
Branch N, from the simulator. The worker has no notification record and no replace ID, so it pushes on every delivery. After event 18 the phone shows two stacked notifications; after event 19's redelivery, three. With a replace ID but still no record: one notification shown, but three alerts. The record removes the repeats it can; the replace ID hides the rest.
Logical duplicates and windows
Some duplicates are not repeats of one request but the same message from different requests: a billing system that raises "your invoice is overdue" twice in a day. The notification loop (step 2.3) handles these with a producer-chosen dedup key (a fingerprint of recipient, template and business key) stored with SET key 1 NX EX <window> in Valkey or Redis:
- The first
SET NXsucceeds and the message goes out; a second within the window finds the key and is dropped. - If the cache is down, fail open: send anyway. A duplicate reminder is better than a lost one.
- Run that cache with
noevictionand alarm on memory: an evicted key silently turns into a duplicate. - Count per-user caps once, at accept, after dedup, so a dropped duplicate doesn't use up the user's allowance.
Email: the send record, and choosing your failure
The email loop (steps 1.6 and round 2's delivery) keeps a send record keyed by the request's Idempotency-Key. It stops our repeats: a client that sends the same request twice finds the record and gets the same answer. It can't stop the provider's: if our call to the email provider times out, the email may or may not have gone. From the simulator, a request sent twice by its client, whose provider call timed out after the provider had accepted it:
| Policy after the unknown | Emails delivered | When to choose it |
|---|---|---|
| Retry the provider call | 2 | Receipts and security codes: a rare duplicate beats a missing one |
| Don't retry | 1 (or 0 if the first call really failed) | Marketing: a rare missing one beats a duplicate |
Third parties with no key
For any outside call that takes no key: look it up by your own reference if the provider offers a search; otherwise accept a rare duplicate, or don't retry after an unknown, and choose per message type. A token that a client library generates for you inside one call covers only that library's retries of that call; if your process restarts and calls again, it is a new token.
Ready means notified
Branch R is the YouTube loop's rule, from step 1.5 ("Is My Video Ready Yet?"). A worker that marks the event DONE, then crashes before sending, leaves a trap: the retry finds DONE and, if it only acknowledges, the buyer is never told. From the simulator: after the crash at 14.50 and the redelivery at 44.50, the phone shows no notification, although m1 was credited. The rule: a detected duplicate checks the notification record and finishes it. "Done" isn't done until the notification is sent.
What to remember from Part 7
- Record the intent with the effect; send; mark sent.
- After an unknown send, re-send with a replace ID, so a notification still shown is replaced, not stacked.
- A detected duplicate must still finish the notification, or nobody is told.
Part 8. Streams: exactly-once state plus idempotent sinks
A stream processor such as Apache Flink advertises exactly-once. It's true, for its state. What it writes to the outside world is another matter: after a crash it replays input and writes some output again. This Part follows a small job through a crash and a restore, and shows what the sink must do.
The job computes merchants' daily totals from the ledger's change stream. It runs alongside the payment story: it reads each change about 4 s after its commit (our assumption), and its watermark (how far it believes event time has got) is simply the newest event time it has read; the change stream carries a heartbeat every second. How real watermarks work is the job of page 07 (Event Time, Watermarks & Checkpoints).
Checkpoints make state exactly-once
A checkpoint saves, together, every operator's state and the input positions that state reflects. A restore loads both, so the job rewinds its state and its input to the same moment and replays from there. Each input event then affects the state exactly once, however many times it is read.
Output is different. Anything the job wrote to the sink after the checkpoint's barrier reached the sink was written by a run that the restore has now erased, and the replay will write it again. So:
Effectively-once = exactly-once state + a sink that makes repeated output harmless.
Our sink writes absolute totals, SET (merchant, day) = total, so writing the same total twice is harmless. Each sink write is also tagged with the id of the last checkpoint barrier the sink has passed (the id Flink hands the sink when it takes its snapshot; after a restore it starts at the restored checkpoint's id), and a sparse index on that tag lists recently written keys. Not the last completed checkpoint: the sink hears of a completion only later, and sometimes never, so writes made between the barrier and that notice would carry an older tag and escape the flush. The ad-click loop added this tag (step 2.4) so a restore can find keys that restored state knows nothing about.
Trace: the job restores
| # | t (s) | Event |
|---|---|---|
| 0.00 | Checkpoint c1 completes. Sink (m1, 09-28) = 12,000 | |
| S1 | 13.00 | Job version v1 reads ev_6, a declined $5.00 payment for merchant m2 (committed at 9.00), and wrongly counts it: sink SET (m2, 09-28) = 500, tagged c1 |
| S2 | 15.00 | The job reads ev_7 (event time 11.00, watermark 11.00). 11.00 + 7,200 = 7,211 > 11, so it isn't late; the dedup state adds ev_7 with a timer at 7,211.00; m1 12,000 → 15,000; sink SET (m1, 09-28) = 15,000, tagged c1 |
| S3 | 20.00 | The job crashes. Its last completed checkpoint is c1 (t = 0; the next was due at 30) |
| S4 | 25.00 | The team restores v2 (declined payments excluded) from c1: state m1 = 12,000, dedup state empty. Full flush first: rewrite every key in restored state (m1 → 12,000, a visible dip); query the tag index for keys tagged c1 or later: m1 and m2; m2 isn't in restored state, so it is deleted. Then replay: ev_6 is not counted (declined; its ID still enters the dedup state, timer 7,209.00), ev_7 → m1 = 15,000, timer 7,211.00 again |
Synthesizing vector architecture diagram...
What to notice: the replay never touches m2, because v2 skips declined payments. Without the flush's delete, the 500 written by v1 would stay in the sink forever.
Snapshot P5, the stream job after S4:
Synthesizing vector architecture diagram...
After S4: the sink matches restored state plus the replay: m1 briefly showed 12,000 during the flush, and m2 is gone.
Idempotent sinks vs two-phase-commit sinks
There are two ways to make output safe to repeat, and they cost different things:
| Absolute, idempotent sink (our job, the ad-click loop) | Two-phase-commit sink (for example Flink's Kafka sink in exactly-once mode) | |
|---|---|---|
| How repeats are made harmless | Replayed writes set the same values again | Output is written inside a transaction that commits only when a checkpoint completes; a restore aborts the uncommitted output |
| When readers see output | Immediately | Only after the next checkpoint completes: up to a checkpoint interval later (Managed Service for Apache Flink's default interval is 60 s) |
| What readers must do | Nothing special, but they can see the flush's brief dip | Read with isolation.level=read_committed; the default read_uncommitted also shows aborted output |
| After a restore | Needs the full flush: keys written after the checkpoint may be stale or orphaned | Needs no flush for output it never committed |
| Settings that can bite | The tag index must be written with every sink write | The transaction timeout must be at least the longest checkpoint plus the longest restart, or, in Flink's words, "data loss may happen when Kafka expires an uncommitted transaction". Flink's Kafka sink sets the producer's transaction.timeout.ms to 1 hour by default, while the broker's transaction.max.timeout.ms defaults to 15 minutes, so one of the two must be changed |
| Where it works | Any store that can do a keyed overwrite and an index | Only sinks with transactions that can be committed later |
The ad-click loop picked the absolute sink because a two-phase-commit sink would have turned its 3-second freshness into the checkpoint interval.
After the restore, m1 briefly shows 12,000 on the merchant's dashboard. Bug or expected?
Branch S, the same restore without the full flush. From the simulator, after the replay the sink holds m1 = 15,000 (correct) and m2 = 500: v1's wrong total for a declined payment, which the replay never touches and nothing will ever rewrite.
textFULL FLUSH(restored checkpoint c) -- after every restore or rescale, before new output 1. for every key in the restored state: SET key = its value, tag = c 2. query the tag index for keys whose tag is c or later -- written after barrier c reached the sink 3. delete every one of them that the restored state doesn't hold 4. replay from c's input positions every sink write carries the id of the last checkpoint BARRIER this sink has passed (after a restore, starting at c), not the last completed checkpoint, whose notice arrives late or not at all; the index on that tag is sparse, so step 2 reads only recently written keys
Dedup state, event time and the acceptance horizon
The change stream can also deliver a change twice (a connector that restarts from an older position). So the job keeps a dedup state: the event IDs it has seen. Two rules make it safe:
- Keep it in the checkpointed state, not in an outside store. From the simulator: with the seen-IDs in an external cache, the cache still holds
ev_7after the restore, the replay ofev_7is dropped as a duplicate, andm1stays at 12,000: the restore under-bills by 3,000. - Expire it by event time, with an acceptance horizon. Each ID gets a timer at its event time + 2 hours (7,200 s); the timer fires when the watermark passes it. And a copy whose event time + 7,200 ≤ watermark is rejected as late (sent to a batch recount). Because the timer and the rejection use the same boundary, every copy that is accepted still finds its original in state, as long as the watermark the rejection uses never goes backwards (rule 3).
- Keep a floor under the watermark. A stream processor like Flink doesn't checkpoint watermarks: after a restore the watermark starts from nothing until the sources send new ones. In that gap a copy that was already too late before the crash would look on time, find its original's ID gone (its timer fired before the checkpoint), and be counted twice. So the dedup operator saves the last watermark it saw in its checkpointed state, and after a restore rejects as late any copy whose event time + 7,200 ≤ the larger of that saved floor and the live watermark.
| # | t (s) | Event |
|---|---|---|
| S5 | 435.20 | The change-stream connector restarted from an older position and re-emitted ev_7 (event time 11.00) at 431.20; the job reads it at 435.20, watermark 431.00. 11.00 + 7,200 = 7,211 > 431, so it isn't late; ev_7 is in dedup state → dropped |
| S6 | 7,215.00 | The watermark reaches 7,211.00: the timer clears ev_7. From now on, including after a restore (rule 3), a copy of ev_7 is rejected as late. The simulator sent one more copy, re-emitted at 7,300.50 and read at 7,304.50 with watermark 7,300.00: 11.00 + 7,200 = 7,211 ≤ 7,300 → rejected as late, not counted |
| Branch T | 10,860.00 | A processing-time TTL instead. No crash at 20; checkpoint c2 completes at 30.00 with ev_7 in dedup state (last touched at 15.00). The job crashes at 59, before c3, and is down from 60 to 10,860 (3 hours) and restores from c2. A 2-hour processing-time TTL on ev_7 expired at 15.00 + 7,200 = 7,215.00, so when the catch-up replay reaches the S5 copy (watermark 431.00), ev_7 is gone: counted again, m1 = 18,000. With the event-time timer (7,211 in event time, and the watermark only at 431), ev_7 is still in state → dropped, m1 = 15,000 |
A processing-time TTL makes dedup depend on how far behind the job is: the longer the outage, the more duplicates. An event-time timer doesn't care. Keep a processing-time TTL only as a much longer guard against state growing without bound.
What to remember from Part 8
- A checkpoint makes state exactly-once; the sink must make repeated output harmless.
- After a restore, rewrite or delete every sink key written since that checkpoint, found through its tag.
- Keep dedup state in the checkpoint, expire it by event time, and reject copies older than the horizon, judged against a watermark floor restored from the checkpoint.
Part 9. Across Regions
Everything so far assumed one Region. With a second Region, every table that remembers a key is replicated, and replication is asynchronous in most designs. That changes what a condition means.
What replication does to a claim
Facts from the DynamoDB documentation for multi-Region eventual consistency (MREC) global tables, the default mode:
- Changes are replicated asynchronously, "typically within a second or less".
- "Conditional writes evaluate the condition expression against the version of the item in the Region": the local replica.
- Conflicting writes are resolved per item by last writer wins.
- Transactions are atomic only in the Region where they run; other Regions may see part of one for a while.
So during the lag, two Regions can both claim the same key, and both succeed. Our story uses a lag of 1.0 s (an assumption).
Trace: two Regions claim the same key (branch F2)
The PSP answers promptly in this branch (no gray failure). A1 claims K1 in R1 at 0.02. At 0.45 the phone's resend A1′ is routed to R2, before R1's claim has replicated (0.02 + 1.0 = 1.02). From the simulator:
Synthesizing vector architecture diagram...
What to notice: one charge, two credits. The derived PSP key protected the one system that dedups in one place, the PSP. Our own tables checked each claim locally, so both Regions ran the whole path.
| Design | Charges | m1 once replication settles | Pushes |
|---|---|---|---|
| Both Regions write; credit entries get Region-local IDs | 1 | 18,000: two credit entries, +6,000 | 2 sent, 1 shown |
| Both Regions write; the credit entry's ID is derived from the event | 1 | 15,000: last writer wins merges the two into one item. Correct here, but by accident | 2 sent, 1 shown |
Both Regions write; the balance is one item updated with ADD, and another $10 payment (pay_9) is credited in R1 at 0.40 | 1 | 15,000, where 16,000 is right: R2's later write of 15,000 wins, and the $10 disappears | 2 sent, 1 shown |
One home Region per key: A1′ is routed to R1 | 1 | 15,000 | 1 |
The first row is the double credit. The other two show that "last writer wins" doesn't rescue you: it can hide a duplicate by accident, and it silently drops any other change to the same item made during the lag.
The fixes:
- A home Region per key. Route every request for a caller (or merchant) to one Region; in the home-Region run, A1′ found
IN_PROGRESSand got409. - Or a claim store that is strongly consistent across Regions: a DynamoDB multi-Region strong consistency (MRSC) table, where "conditional writes always evaluate the condition expression against the latest version of an item". MRSC runs in exactly three Regions (three replicas, or two plus a witness), supports no transactions and no TTL, and a write that conflicts with one in progress in another Region fails with
ReplicatedWriteConflictException. - Or a single-writer ledger: one Region takes all ledger writes.
Trace: a Region fails in the middle of replication (branch F)
The main line, now with a standby Region R2 that is 1.0 s behind. From the simulator:
| t (s) | What happens |
|---|---|
| 11.00 | R1's resolver commits pay_7 SUCCEEDED (ch_1) and the outbox row ev_7 = pay_7:SUCCEEDED. With 1.0 s of lag, R2 would have it at 12.00 |
| 11.20 | R1's relay publishes ev_7 to R1's queue |
| 11.30 | R1 goes dark. R2 holds R1's writes up to 10.30: pay_7 UNKNOWN, and the key record's pointer answer |
| 60.00 | R2 is promoted (an assumed promotion time; detection and promotion are page 09's job). R1's workers stay stopped |
| 61.00 | R2's resolver finds pay_7 UNKNOWN and replays pay_7:charge with identical parameters. It is 61 s after the charge, well inside the PSP's 24 hours: the PSP returns the saved ch_1. pay_7 → SUCCEEDED; R2 re-creates the event with the same derived ID pay_7:SUCCEEDED |
| 61.50 | R2's worker: inbox insert, 1 row; credit m1 → 15,000; one push |
| later | R1 rejoins as a replica. Its queue still holds ev_7; drained against the home Region's inbox, it inserts 0 rows and is dropped. Reconciling R1's lost tail by ID: R1 said pay_7 SUCCEEDED with ch_1, and so does R2. One payment, one credit |
Synthesizing vector architecture diagram...
What to notice: R2 never saw R1's last second, yet it produced the same charge ID and the same event ID, because both are derived. That is what lets the PSP and the inbox recognise the replay.
Two counterfactuals, also from the simulator:
- A per-Region attempt key. If
R2replays withpay_7:R2:1, the PSP has never seen that key and charges again:ch_2, $60.00 in total. - A random event ID in
R1. IfR1's outbox row had carried a random ID whileR2derivespay_7:SUCCEEDED,R1's leftover copy doesn't match anything in the home inbox, so it is credited again:m1= 18,000. This is why the event ID is derived in any design with a Region failover (Part 2); a single-Region design may store a random one.
R2 was promoted with pay_7 still UNKNOWN. Its resolver replays the charge. Why is that safe, and what key would have made it unsafe?
Where a claim can live
| Single-Region table | MREC global table | MRSC global table | Single-writer ledger (for example an Aurora global database) | |
|---|---|---|---|---|
| Where a claim is checked | In one place | On the local replica | Against the latest version | In the primary Region's writer |
| Can two Regions both claim? | No | Yes, during the lag | No: a conflicting write fails with ReplicatedWriteConflictException | No |
| Transactions | Yes | Atomic only in their own Region | None | Full SQL transactions |
| TTL | Yes | Yes | None: expire in the read path and clean up yourself | Not applicable |
| Claim latency from another Region | A cross-Region round trip | Local | Includes a synchronous write to at least one other Region | A cross-Region round trip to the writer |
| When its Region fails | Claims stop until you recover or promote | The other Regions keep going; writes not yet replicated arrive only when the failed Region recovers | AWS documents a recovery point objective of zero: no acknowledged write is lost | Promote a secondary (Aurora replicates "with latency typically under a second"); resume UNKNOWNs with the same keys and reconcile the tail by ID, as in branch F |
Regional dedup windows don't cross Regions. A FIFO queue's dedup IDs, a Valkey fingerprint cache and a Kafka producer session all live in one Region. After a failover, the new Region's copies of them start empty.
For quorums, lag and read-your-writes, see page 06 (Replication, Quorums & Read-Your-Writes); for detection, promotion and routing, page 09 (Multi-Region Failover). Meanwhile, Primitive #24: Disaster recovery and multi-Region covers the broad picture.
What to remember from Part 9
- Async multi-Region tables check a claim locally: during the lag, two Regions can both win and both commit.
- Give each key one home Region. Derived IDs only make outside systems recognise a replay; your own tables still need one writer.
- After a failover, resume
UNKNOWNs with the same keys, and reconcile the failed Region's tail by ID.
Part 10. Do you need it? Compared on equal terms
Idempotency keys cost storage, a write on the hot path and discipline in every client. Before paying that, it is fair to ask whether something people already know (a lock, a distributed transaction, a saga, an exactly-once broker) would do the job. Each gets the same test: the gray failure of events 3 to 9, where the PSP charges and its reply is lost.
The payment six ways
From the simulator, the same failure under each approach (the last row is this page's design):
| Approach | What happens | What the buyer is charged, and told |
|---|---|---|
| Never retry | The server's PSP call times out; with no UNKNOWN state it marks the payment failed. The phone times out and gives up | $30 charged, told "payment failed" |
| Check, then act ("is there a payment for order 42? if not, create one") | The first attempt writes its row only after the PSP answers, which never happens. A1′ at 0.45 checks, finds nothing, charges again | $60 charged, told "paid" |
| A lock or lease on the order | A1′ is refused while A1 holds the lock. A1 times out at 5.02, releases, marks failed. A2 at 10.00 finds the lock free and nothing remembered, and charges | $60 charged, told "paid" |
| Two-phase commit of the ledger and the outbox | The PSP has no "prepare", so the charge happens outside the transaction. The transaction aborts at its timeout; A2 runs as a new transaction and charges | $60 charged, told "paid" |
| A saga with idempotent steps, named after the client's key | A2 finds the same saga running; its charge step retries with pay_7:charge and the PSP replays ch_1 | $30 charged, told "confirming", then "paid" |
| Exactly-once in the broker | Each request becomes a message; the consumer charges. Kafka's transactions cover writes to the log; the PSP call and the lost reply are outside. A2 is a new message | $60 charged, told "paid" |
| Idempotency key + durable guard | A1′ gets 409; A2 gets the pointer answer; the resolver settles pay_7; the push says "confirmed"; A3 gets 200 from the ledger row | $30 charged, told "confirming", then "paid" |
The saga row works for the same reason the last row does: it is named after the client's key, and each step carries a derived key. Take the keys away and it behaves like the two-phase-commit row.
The comparison, on equal terms
| Covers | Needs | Costs | Fails when | Answers a client whose reply was lost? | |
|---|---|---|---|---|---|
| Never retry | No duplicates at all | Nothing | Every unknown becomes a wrong "failed" | Any timeout | No |
| Check, then act | Repeats in sequence | Correct only inside one serializable transaction in one database (Primitive #21: Isolation levels); across a database and a PSP it races | A read before each write | Two copies run at once; or the first copy hasn't written its row yet | Only if the first attempt left a record |
| Lock or lease | Two copies at the same time, while the lock is held | A lock service, expiry and a fence at the resource (Leases, Fencing Tokens & Distributed Locks) | A lock round trip; a gap after a crash | The same request comes back after the holder finished | No |
| Two-phase commit (XA) | Several stores that commit together | A durable coordinator log; participants that support prepare | Locks held across two round trips; a coordinator crash after prepare blocks the participants | Any participant can't prepare (a PSP, an email provider) | No: the client's retry is a new transaction |
| Saga + idempotent steps | Multi-step work across services, with undo | Idempotent steps (keys), derived compensation keys, a durable orchestrator | No isolation; compensation logic for every step | A step without a key repeats its effect | Yes, if the saga is named by the client's key |
| Broker exactly-once | Read, process and write inside the log | Transactional IDs; read_committed readers | Output visible only after commit | Effects outside the log; a producer's new record from an outbox re-send | No |
| Key + durable guard | Every repeat that carries the key, for as long as the payment row lives | Client discipline; a key store; a unique key and fingerprint on the durable record | Storage, one conditional write before the external call, stored answers | A client that loses its key (a reinstall: Part 4's per-device last ID) | Yes: it replays the answer or renders the payment's state |
Why not wrap the ledger write and the PSP charge in one two-phase commit?
What a saga gives up
A saga runs a series of local transactions, each with a compensation that undoes it (a refund undoes a charge). Each step commits on its own, so a saga gives up isolation: other readers can see the intermediate states. A merchant's dashboard can show a payment as charged, then refunded; a stock check can see an item reserved by an order that is about to be cancelled. Designs handle this with semantic counter-measures: a PENDING state that readers treat as not final, updates that give the same result in any order, and re-reading a record before acting on it. Its compensations are ordinary steps too: each has its own derived key (pay_7:cancel:refund), and an unknown one is retried, not abandoned (Part 4).
Why teams choose a saga over two-phase commit across services (the second question of the drill The Flight Booking That Charged Without a Seat): no locks are held across services or across a slow partner API; a crashed coordinator doesn't leave other services blocked, because each step already committed and the orchestrator resumes from its durable log; each service stays available on its own; and each step pays one local commit, not two rounds. The price is the lost isolation above and the compensation logic. For how orchestration works, see Primitive #10: Two-phase commit and saga orchestration.
What to remember from Part 10
- Two-phase commit makes several stores commit together; it never tells a client whose reply was lost what happened.
- A lock stops two runs at once, not the same run twice; only a remembered key stops the repeat.
- Every hop needs its own ID and its own guard; one exactly-once hop doesn't make the chain exactly-once.
Part 11. End to end: from the tap to the notification
"We made the API idempotent" protects one hop out of nine. A repeat can be born anywhere along the path, and each place needs its own ID and its own memory. This Part lists them all.
Every place a repeat is born
Synthesizing vector architecture diagram...
What to notice: every box except the load balancer, the ledger and the push provider can repeat something. We don't draw a retry at the load balancer: this page relies only on its arrival stamp, not on any retry behaviour.
The hop table
| Hop | Who repeats | The ID | Where it is remembered | For how long | What decides after that |
|---|---|---|---|---|---|
| Phone → API | The phone's HTTP library may resend after a connection reset; the app retries on its timer and after restarts | K1 | Key store idem, and the payments row's unique key | Key store: 1 hour after completion. Row: as long as the payment | The row's UNIQUE (u1, K1) and fingerprint |
| Load balancer | Nothing we rely on | Its arrival stamp | On the request | That request | The server's deadline |
| API's own AWS calls (for example the key-store claim) | The AWS SDK: 3 attempts in total by default (4 for DynamoDB under AWS's updated retry behaviour; older defaults vary) | The conditional put itself | The key-store item | Its lifetime | A retried claim that already succeeded fails its own condition: read the item and check the owner before calling it a conflict |
| API → PSP | The API before its deadline; the resolver; another Region's resolver after failover | pay_7:charge | The PSP's key store | At least 24 hours | Look the payment up by our reference |
| PSP → our webhook | The PSP, for up to three days | evt_1 | webhook_events | Longer than the sender retries | The payment's state machine (compare-and-set) |
| Ledger → relay → queue | The relay re-sends rows it didn't delete | ev_7 = pay_7:SUCCEEDED | The queue's dedup window | 5 minutes | The worker's inbox |
| Queue → worker | The queue, after the visibility timeout | ev_7 | Inbox processed | Longer than the longest relay outage | Not applicable: the inbox is the guard |
| Worker → push provider | The worker, when notif is still PENDING | Intent (ev_7, push); replace ID pay_7 | notif rows; the phone's notification drawer | Rows: like the inbox; drawer: while shown | A duplicate is possible; the replace ID hides it only while the first is shown |
| Ledger → stream job | The change-stream connector, and every restore | ev_7 | The job's checkpointed dedup state | 2 hours of event time | Rejected as late, recounted in batch |
One tap, one sequence
A1 end to end, with the ID each hop carries:
Synthesizing vector architecture diagram...
What to notice: the IDs change name at each hop but never at random: K1 names pay_7, pay_7 names the PSP key, the event and the replace ID. That chain is the design.
What to remember from Part 11
- Every hop has something that repeats; list them before designing.
- Each hop needs an ID, a store that remembers it with the effect, and a lifetime longer than that hop's repeats.
- The ID chain is the design: key → payment → PSP key; event ID → inbox → notification intent.
Part 12. On AWS
The mechanism is the same on AWS. What changes is which service plays the durable guard, which ones offer short windows, and which only look like dedup.
Managed services that use it
| Service | What it uses | What the documentation says |
|---|---|---|
| Amazon DynamoDB | Conditional writes for key claims, inboxes and state moves; TransactWriteItems with a ClientRequestToken | A transaction holds up to 100 actions. A client request token "is valid for 10 minutes after the first request that uses it is completed", then a repeat is a new request; a changed request inside the window gets IdempotentParameterMismatch; tokens are 1 to 36 characters; repeats have no further side effects but may report different consumed capacity. TTL deletes expired items "within a few days". MREC global tables check conditions locally, resolve conflicts by last writer wins, and keep transactions atomic only in their own Region; MRSC tables run in exactly three Regions with no transactions and no TTL, and conflicting writes can fail with ReplicatedWriteConflictException |
| Amazon SQS | FIFO MessageDeduplicationId or content-based dedup; redelivery after the visibility timeout | 5-minute dedup window; a repeat is accepted but not delivered; the ID is tracked even after the message is received and deleted; content-based dedup is a SHA-256 of the body, not the attributes; high-throughput mode scopes dedup to the message group. Visibility timeout: 30 s by default, at most 12 hours from first receipt |
| Amazon MSK | Kafka's idempotent producer and transactions | Producer defaults: enable.idempotence true, acks all, at most 5 in-flight requests, delivery.timeout.ms 120,000, transaction.timeout.ms 60,000. Consumer isolation.level defaults to read_uncommitted. Broker transaction.max.timeout.ms defaults to 15 minutes |
| Amazon Managed Service for Apache Flink | Checkpoints (exactly-once state); Flink's Kafka sink in exactly-once mode | Default checkpointing: enabled, every 60,000 ms, at least 5,000 ms between checkpoints. Flink's Kafka sink commits transactions on checkpoints, needs a unique transactional ID prefix, and its transaction timeout must cover the longest checkpoint plus restart, or "data loss may happen"; the sink's own default transaction timeout is 1 hour |
| AWS Lambda with Powertools for AWS Lambda | The idempotency utility, over a DynamoDB table | Records expire after 3,600 s by default (expires_after_seconds); states INPROGRESS, COMPLETE, EXPIRED; a concurrent duplicate gets IdempotencyAlreadyInProgressError; payload validation raises IdempotencyValidationError. An unhandled exception deletes the record, so a function that called something and then failed will call it again: that callee must be idempotent too |
| AWS Step Functions | Standard workflows; StartExecution idempotent by name | Standard: same name and input while running returns the original response; after it closes, ExecutionAlreadyExists; the name can be reused 90 days after close. Not idempotent for Express workflows |
| Amazon EC2 (the API) | ClientToken on RunInstances and many other actions | Up to 64 ASCII characters; a retry with the same token and parameters "succeeds without performing any further actions" (the result may show newer status); different parameters give IdempotentParameterMismatch. How long a token is remembered isn't stated in the documentation |
| AWS Elemental MediaConvert | CreateJob client request token | "If you reuse a client request token within one minute of a successful request, the API returns the job details of the original request" |
| Amazon S3 | Fixed output keys written as whole objects; conditional writes | If-None-Match prevents overwriting an existing key; If-Match checks the ETag before writing |
| Amazon Aurora (PostgreSQL) | The durable guard: a unique key kept forever; INSERT … ON CONFLICT DO NOTHING; the key row in the payment's transaction | A second insert of a key that an uncommitted transaction has inserted waits to see whether that transaction commits, then conflicts (or goes ahead if it rolled back). Global databases replicate to secondary Regions "with latency typically under a second" |
Running it yourself
| Option | What it is | Sizing and notes |
|---|---|---|
| Amazon ElastiCache (Valkey): a building block, not a service that dedups for you | SET key value NX EX <window> in front of a durable guard (the notification loop's fingerprints) | Replicas are updated asynchronously, so after a failover "there may be some data loss": recent keys can be forgotten, which shows up as duplicates. Fail open, run with noeviction, alarm on memory |
| A key store at the payments loop's round-2 scale | A DynamoDB table of key records in front of Aurora's unique key | 20 million payments a day × 7 days × about 500 bytes = 70 GB, plus more because TTL deletes lag. Each payment adds two conditional writes before the PSP call (the claim and the PENDING row) |
Shared limits
- The PSP's rate limit is shared by live charges and the resolver's replays. Stripe documents default live-mode limits of 100 requests a second per account and 25 a second per endpoint unless noted, and asks high-volume users to request increases in advance. Whatever your limit, give the resolver its own cap, so replays after an outage don't starve new payments.
- The durable guard's unique index is written by every create and shares the ledger's single writer with everything else.
- One merchant's balance item or row takes every credit for that merchant (Part 6).
- The FIFO message group serialises a merchant's events.
- The dedup cache's memory holds every fingerprint in its window; eviction turns into silent duplicates.
- The full flush shares the sink's write capacity with live output.
Look-alikes that are not this mechanism
| Look-alike | Why it looks like it | Why it isn't |
|---|---|---|
| SQS FIFO "exactly-once processing" | The name | Dedup of sends for 5 minutes; consumers still see redeliveries after a visibility timeout and must dedup themselves. SNS FIFO's publish dedup has the same 5-minute window |
| Kinesis Data Streams | An ordered log, like Kafka | AWS names two causes of duplicate records, producer retries and consumer retries, and says the application "must anticipate and appropriately handle processing individual records multiple times" |
| DynamoDB TTL | "The key expires" | Deletes within a few days; it is cleanup, never the window |
| DynamoDB global tables (MREC) | Replicated conditional writes | Conditions run on the local replica; last writer wins; two Regions can both claim and both commit (Part 9) |
| Step Functions Express | "Step Functions is exactly-once" | StartExecution isn't idempotent for Express workflows |
| CloudFront request collapsing and Origin Shield | "Duplicate requests become one" | They merge simultaneous requests for the same cached object into fewer origin fetches: reads at the same moment, not a write repeated later |
| Content-addressed dedup (hash-keyed objects, a file-sync chunk store) | "Dedup" | Stores identical bytes once; says nothing about repeated requests |
| Amazon SES v2 and End User Messaging SMS v2 | Messaging services on the path | SendEmail and SendTextMessage take no idempotency token; repeats must be prevented, or accepted, by us (Part 7) |
What to remember from Part 12
- DynamoDB conditional writes and a unique key in Aurora are the durable guards; tokens with 1-to-10-minute windows are accelerators.
- SQS FIFO drops repeated sends for 5 minutes; consumers still dedup.
- Managed Flink gives exactly-once state; the sink's idempotence is yours.
Part 13. What you've learned
Back to the tap
The buyer tapped Pay once, and the network lost the answer at almost every hop. The simulator's final count: one charge (ch_1), one credit to m1 (12,000 → 15,000), one notification on the phone. Here is how each piece did its job:
- The app made
K1at the tap and reused it on every retry (Part 2). The server recorded the intent before the PSP call: aPENDINGrow, unique on(u1, K1)(Part 3, snapshot P1). - The PSP's reply was lost. At the deadline the payment became
UNKNOWN, not failed, and the key record stored a pointer, so the retry at 10.00 heard "processing" rendered from the payment itself (Part 3, snapshot P2). - The resolver replayed the derived PSP key and got the saved
ch_1; a compare-and-set made the late webhook harmless (Part 5). - The relay sent
ev_7three times. The queue's window dropped the second; the worker's inbox dropped the third (Part 6, snapshot P3). - A worker died after pushing. The notification record made the next worker finish the job, and the replace ID made the second push replace the first (Parts 6 and 7, snapshot P4).
- The stream job crashed and restored. Exactly-once state, an absolute sink and the full flush left the totals right (Part 8, snapshot P5).
- The buyer came back after two hours. The key store had forgotten
K1; the durable guard answered (Part 4, snapshot P6). - In the branches: a deadline and a void stopped the late original that opened this page; a home Region stopped the double credit in F2; derived IDs made
R2's replay safe in F (Parts 4 and 9).
The whole story, event by event
Payment timeline:
| # | t (s) | Event |
|---|---|---|
| 1 | 0.00 | Tap: K1 created, journal SENDING; A1 sent over Wi-Fi |
| 2 | 0.02 | Load balancer stamp 0.02; api-1 claims K1 (IN_PROGRESS, H1, until 5.02); inserts pay_7 PENDING |
| 3 | 0.05 | api-1 calls the PSP with pay_7:charge |
| 4 | 0.30 | The PSP charges ch_1; its reply is delayed |
| 5 | 0.45 | Wi-Fi dies; the HTTP library resends A1′ to api-2: 409, Retry-After: 1 |
| 6 | 5.02 | Deadline: pay_7 PENDING → UNKNOWN; idem COMPLETED with a pointer, expires_at 3,605.02 |
| 7 | 5.03 | The 202 has nowhere to go: A1's connection is dead |
| 8 | 5.60 | The PSP's reply arrives on a closed connection: lost |
| 9 | 10.00 | Timer: A2 → 202 PROCESSING rendered from pay_7; journal UNKNOWN |
| 10 | 10.60 | The buyer closes the app |
| 11 | 11.00 | Resolver (UNKNOWN for 5.98 s) replays pay_7:charge → saved ch_1; pay_7 SUCCEEDED; outbox ev_7; final answer stored |
| 12 | 11.20 | Relay sends ev_7: accepted, window to 311.20; crashes before deleting the row |
| 13 | 11.50 | Webhook evt_1: new; compare-and-set finds SUCCEEDED, no change |
| 14 | 14.00 | Relay re-sends ev_7, 2.80 s into the window: acknowledged, not enqueued; crashes again |
| 15 | 14.50 | w1: inbox 1 row; m1 12,000 → 15,000; notif PENDING |
| 16 | 14.60 | w1 pushes, replace ID pay_7 |
| 17 | 14.65 | w1 crashes before SENT and before deleting the message |
| 18 | 44.50 | w2 gets ev_7: inbox 0 rows; notif PENDING → pushes again at 44.60 (replaces); SENT; deletes the message |
| 19 | 431.20 | Relay re-sends the row, 120 s after the window: delivered; inbox 0 rows, notif SENT; the row is deleted |
| 20 | 7,200.00 | A3: idem expired (3,605.02 < 7,200); the pay_7 insert hits UNIQUE (u1, K1): 200 {pay_7, SUCCEEDED} with a replay marker |
Stream job, parallel timeline:
| # | t (s) | Event |
|---|---|---|
| S1 | 13.00 | v1 counts declined ev_6: sink (m2) = 500, tag c1 |
| S2 | 15.00 | ev_7 (event time 11.00): timer 7,211.00; sink (m1) = 15,000, tag c1 |
| S3 | 20.00 | The job crashes; last checkpoint c1 |
| S4 | 25.00 | Restore v2 from c1; flush: m1 → 12,000, m2 deleted; replay: m1 = 15,000 |
| S5 | 435.20 | A re-emitted ev_7, watermark 431.00: in dedup state → dropped |
| S6 | 7,215.00 | Watermark reaches 7,211.00: the timer clears ev_7; later copies are rejected as late |
The cheat card
| Topic | Remember |
|---|---|
| Three outcomes | Done, not done, unknown. A timeout is unknown; an error on a retry may be a success in disguise |
| The key | Made by the client, once per intent, stored before the first send, reused on every retry; no personal data |
| Scope and fingerprint | Endpoint + caller + key; a hash of the request's meaning. Same key, other meaning → 422 |
| Order of work | Claim the key → insert the PENDING row (unique on caller + key) → external call → compare-and-set the state → complete the key only if you still own it |
| Stored answers | Final answers whole, rejections included; a non-final answer as a pointer to the resource |
| Concurrent repeat | 409 + Retry-After (two stores); in one PostgreSQL transaction, the second insert waits, then conflicts |
| Takeover | After in_progress_until + a margin larger than the clock offset; conditional on the owner; check the durable record first |
| Windows | Accelerators. Correctness lives on a record that exists as long as the effect. Expired = expires_at < now in the read path; TTL is cleanup |
| Late original | Server deadline from the load balancer's stamp + the wait before the stamp < client timeout; a "did it happen?" call voids an absent key in both stores |
| IDs | Derive what someone else may re-create (payment ID, PSP key, event IDs with a Region failover); a stored ID may be random |
| Writes | Absolute values, set-if-larger, insert-if-absent, compare-and-set, versions from one counter; never a phone timestamp |
| Queues and logs | FIFO dedup: 5 minutes. Idempotent producer: one session. Transactions: read_committed. The consumer's inbox catches the rest |
| Notifications | Intent row with the effect → send → mark SENT; re-send PENDING with a replace ID (apns-collapse-id, Android tag; not FCM collapse_key) |
| Streams | Exactly-once state + absolute sink; full flush after every restore via a checkpoint tag; dedup state in the checkpoint, expired by event time, with an acceptance horizon and a watermark floor across restores |
| Regions | Async tables check claims locally: a home Region per key, or MRSC, or one writer; resume UNKNOWNs with the same keys; reconcile the tail by ID |
Failure checklist
- Does every client make its key before the first send, store it, and reuse it on every retry?
- Is the durable record of the intent written before any external call, with a unique key that lives as long as the effect?
- Does the read path treat an expired key record as absent, and does the durable guard answer after the window?
- Is the server's deadline, counted from the load balancer's stamp, shorter than the client's timeout minus the wait before the stamp? Does "did it happen?" void an absent key?
- Is every condition on the record that never moves or expires (or leaves a tombstone)?
- Is every write from an at-least-once consumer absolute, versioned or guarded by an inbox, never an increment?
- Does every re-send (relay, resolver, standby Region, re-run parser) carry the original or derived ID, never a fresh one?
- Does anyone rely on a broker's or producer's dedup as the only guard?
- Does a detected duplicate finish the steps the first attempt may not have finished (re-publish, re-send the notification)?
- Is every notification recorded as an intent with the effect, and sent with a replace ID?
- After a stream-job restore, does a full flush rewrite state and delete keys tagged with the checkpoint's id or later (tagged by barrier, not completion)? Is dedup state checkpointed and expired by event time?
- Does the late check use a watermark floor restored from the checkpoint?
- In more than one Region, does each key have one home Region (or a strongly consistent claim store), and does a failover resume unknowns with the same keys?
Think-first drills
Drill 1. A client times out after 8 s, asks with a plain lookup, and may then make a new key. A request can wait up to 12 s in the load balancer and the server's thread pool before the server stamps it; your server has no deadline; the key store's in-progress lease is 10 s. What can go wrong, and which numbers do you set?
Drill 2. A FIFO queue drops repeats for 5 minutes; your relay can be down for 20 minutes and re-sends on recovery. Where does the durable guard go, and what ID must the re-send carry?
Drill 3. A stream job writes ADD 1 per click to a table. After a restore, the totals are too high. Give the two changes that fix it, and say why a restore still needs one more step.
Interview questions
| Question | Model answer |
|---|---|
Design idempotency for POST /payments. | The client makes a key per intent and stores it before the first send. Scope it to endpoint + caller; store a fingerprint of the request's meaning. Claim the key with a conditional write (IN_PROGRESS, owner, in_progress_until = arrival + deadline), then insert the payment PENDING with UNIQUE (caller, key) before calling the PSP with a key derived from the payment ID. Complete the key only if you still own it: final answers stored whole, rejections included; non-final answers as a pointer. Repeats: replay, 409 while in progress, 422 on a different fingerprint. |
| The PSP call timed out. Did it charge, and what do you do? | Unknown. Mark the payment UNKNOWN, tell the client "processing", and resolve: replay with the same key and identical parameters inside the PSP's window (it returns the saved result, even a saved error), or look the payment up by our reference after it. Keep retrying with backoff; after a budget, alert a human; reconciliation has the last word. Never fail over to another PSP while the outcome is unknown. |
| How long do you keep idempotency keys, and what happens after? | The fast key store keeps them for a window (hours to days) as a speed-up. Correctness comes from a durable record that lives as long as the effect: the payment row's unique key and fingerprint. After the window, a retry's insert hits that unique key and the server answers with the payment's current state. Expiry is checked in the read path; TTL is only cleanup. |
| "Kafka has exactly-once." What does it cover, and what doesn't it? | The idempotent producer drops its own retries within one session; a restarted producer without a transactional ID is a stranger. Transactions make a read-process-write cycle inside Kafka atomic, and readers need read_committed. Not covered: effects outside Kafka (a PSP charge, a push, a database write), a client whose reply was lost, and an outbox re-send, which is a new record. Consumers still need an inbox or idempotent writes. |
| How do you avoid sending a push notification twice? | You can't make it impossible: providers take no key. Record a notification intent in the same transaction as the effect, send, mark SENT; a redelivery that finds PENDING sends again. Use a replace ID (apns-collapse-id, the Android notification tag) so a repeat replaces a notification still shown. FCM's collapse_key only affects messages not yet delivered. |
| Why not two-phase commit? | The PSP can't prepare; a coordinator crash after prepare blocks everyone holding locks; and 2PC never answers a client whose reply was lost: its retry is a new transaction. Use idempotent steps with derived keys, a saga for multi-step undo (accepting weaker isolation), and a durable record of each step's state. |
Where to go next
- Write-Ahead Log, fsync & Group Commit: why an error or a timeout leaves a write's outcome unknown.
- Leases, Fencing Tokens & Distributed Locks: stopping a zombie holder with a fencing token, and the clock-and-margin rule a takeover uses.
- Change Streams & the Transactional Outbox: how an event leaves the database, relays vs change data capture, ordering, and apply-by-version consumers. Also Primitive #12: Change data capture and the outbox pattern.
- The Replication, Quorums & Read-Your-Writes (page 06), Event Time, Watermarks & Checkpoints (page 07) and Multi-Region Failover (page 09) loop primitives, all coming. Meanwhile, Primitive #24: Disaster recovery and multi-Region.
- Primitive #05: Message queues vs event streams and Primitive #10: Two-phase commit and saga orchestration.
- Drill: The Flight Booking That Charged Without a Seat, answered in Parts 4 and 10.
- Loops:
- Payments (steps 1.1, 1.2, 1.4, 2.1, 2.4, 3.5, and round 2's key store), digital wallet (1.4, 2.2, 2.5), hotel reservations (1.5, 1.6, 3.5), transactional outbox and ledger (1.3, 1.5, 3.4, 3.5), stock exchange (3.4).
- Message queue (1.2, 2.2, 3.2), key-value store (round 2, 3.2), rate limiter (3.2), URL shortener (3.2).
- Chat (1.3, 1.4, 2.2), leaderboard (1.3, 3.2), notifications (1.3, 2.3, 2.4, 2.6, 3.4).
- Job scheduler (1.3, 1.4, 2.4, 3.3, 3.4), S3-like storage (3.1, 3.5), web crawler (1.1).
- Proximity (2.1), nearby friends (2.6), ride-sharing (round 1, 3.6).
- YouTube (1.5, 2.2, 2.5), Drive (3.5, 3.6), ad-click aggregation (1.3, 2.4, 2.6, 3.3), email (1.3, 1.6, round 2, 3.4).
- Mobile news feed (2.2), mobile chat (1.2, 2.2), mobile stock trading (1.2).