You're buying a pair of headphones for $49.99. You press Pay, and behind the button the shop's code has two jobs to do. It has to charge your card, which the payments system takes care of, and it has to record order number 2, which the orders system takes care of. If both jobs finish, the headphones arrive at your door.
The two systems are separate programs with separate databases, and that's where it gets hard. A database can make several changes happen all together or not at all by wrapping them in a transaction: you start one, make the changes, and then either commit them or roll them back, and the database never lets anyone see a half-finished result. But a transaction belongs to one database. Nothing covers the payments database and the orders database together, so if the shop's server loses power after the charge and before the order is saved, you've paid for headphones that, as far as the shop's records go, you never ordered.
The question for this chapter is how two systems that fail independently can agree on whether order 2 happened. We'll follow that one checkout the whole way. First we'll crash it halfway and look at the damage. Then we'll make two databases commit together with two-phase commit and watch where that breaks down, and after that we'll look at what real services use instead: sagas, the outbox and idempotency keys. The chapter ends with a careful look at what Kafka means when it says "exactly-once".
01One checkout, two systems
1.1A charge with no order
We can reproduce the failure in a short script. It uses two separate in-memory SQLite databases (SQLite is a small database that runs inside your own program) as stand-ins for the two systems, one for orders and one for payments. The function checkout does what the shop's code does: charge the card first, then save the order. Python's with connection: block runs its contents as one transaction and commits when the block ends, so each write below is a properly committed transaction in its own database. Checkout 1 runs normally. Checkout 2 is told to crash in between, after the charge has committed and before the order is written, which stands in for a power cut or a killed process.
Checkout 2 crashes after charging the card and before saving the order. What will the two databases hold afterwards?
Save the script as dual.py and run it.
import sqlite3
orders = sqlite3.connect(":memory:") # one system
payments = sqlite3.connect(":memory:") # a second, independent system
orders.execute("CREATE TABLE orders (id INTEGER PRIMARY KEY, item TEXT)")
payments.execute("CREATE TABLE charges (order_id INTEGER, amount INT)")
def checkout(order_id, crash_between=False):
with payments: # charge the card and commit in system 2
payments.execute("INSERT INTO charges VALUES (?, 4999)", (order_id,))
if crash_between:
raise RuntimeError("process crashed before saving the order")
with orders: # then save the order and commit in system 1
orders.execute("INSERT INTO orders VALUES (?, 'headphones')", (order_id,))
checkout(1)
try: checkout(2, crash_between=True)
except RuntimeError as e: print("checkout 2 failed:", e)
print("orders :", orders.execute("SELECT id FROM orders").fetchall())
print("charges :", payments.execute("SELECT order_id FROM charges").fetchall())checkout 2 failed: process crashed before saving the order
orders : [(1,)]
charges : [(1,), (2,)]Read the last two lines. orders has one row, for order 1. charges has two rows, for orders 1 and 2. The amount in the script is 4999, the price in cents, so customer 2 has been billed $49.99 for an order the shop has no record of. Here is the same run drawn out, with the two systems as separate places:
1.2Why a transaction doesn't save you
Each with block did its job, because each one was atomic inside its own database. The trouble is that no single commit covers both. To see why, look at what a commit physically is. Every database keeps a log, a file it only ever appends to, and a transaction commits at the moment one record saying "this transaction committed" is safely on disk in that log. One record in one log is what the whole guarantee rests on. It can cover everything in that database and nothing outside it: no record in the payments log can say "and the orders database has this too".
Could we reorder the writes, saving the order first and charging second? That only moves the failure: a crash between the two now leaves an order nobody paid for. The same shape appears wherever code changes a database and then tells another system, and the next section shows another case, with a second and subtler failure hiding behind it.
02Why retrying doesn't fix it
2.1The dual write
The checkout has this problem in a second place. After the order is saved, the warehouse has to hear about it so that someone starts packing the headphones. The shop does that by publishing an OrderPlaced message to Kafka, a system that stores streams of messages in order. Each stream is called a topic, and other services read it at their own pace. Most first versions write to one system and then the other (the order goes into Postgres, a popular open-source database, in this example):
with db.transaction():
db.execute("INSERT INTO orders ...")
kafka.send("orders", order_placed_event) # a second, separate writeThis is called a dual write: two writes to two systems, each atomic on its own, with nothing tying them together. There are four ways it can end, and only one of them is correct.
| Postgres | Kafka | Result |
|---|---|---|
| Committed | Sent | Correct |
| Committed | Process died before send | Order exists; the warehouse never hears about it |
| Committed | send timed out, then succeeded on retry | Correct, or two events if the first attempt landed too |
| Rolled back | Sent (if you publish first) | The warehouse packs an order that doesn't exist |
?Why not just swap the order, or retry until it works?
Swapping the order swaps which row of the table you get, and doesn't change whether you get one. Retrying doesn't cover the case where the process itself dies between the two writes, because then nobody is left to retry. And as section 1.2 said, each system has its own log, and a commit is a single record in one of them, so there's no record Postgres could write that says "and Kafka also has this".
The third row of that table needs a closer look, because it hides a different kind of problem.
2.2The ambiguous timeout
When a call across a network times out, meaning no answer came back within the time you were willing to wait, you don't know whether the call happened. Here is the payment call for order 2, with the answer lost on the way back.
A timeout has three possible meanings: the request never arrived, it arrived and failed, or it arrived and succeeded. From the client's side all three look the same. Every technique in this chapter is a way of making that ambiguity safe, and they come in three kinds: have a coordinator that can find out what happened, have a way to undo it, or have a way to repeat it harmlessly. We'll build them in that order, starting with the one that sounds most natural, which is to get both systems to commit together.
03Making two systems commit together
3.1The protocol
If we want real atomicity across two databases, we need one party to make a single decision and every database to follow it. That protocol is two-phase commit, usually shortened to 2PC. Jim Gray described it in his 1978 Notes on Data Base Operating Systems, and it sits under most of the "distributed transaction" features you'll meet. It has two roles. The coordinator runs the transaction and makes the decision. The participants are the databases that hold the data, which for order 2 are Orders and Payments. For our checkout the coordinator would be the checkout code itself, or a separate transaction manager it hands the job to. The trick is to split "commit" into a promise and a decision.
Three words first, because the protocol leans on them. A lock is how a database stops anyone else from changing a row while a transaction is working on it, held until that transaction commits or rolls back. Durable means written to disk, so that it survives a crash and a restart. And to abort a transaction is to roll it back.
?Why does it need two phases?
Because a participant might be unable to commit: a constraint fails, it runs out of disk, it crashes. Phase one finds that out while it's still safe to abort everywhere. Once every participant has promised, nothing but the coordinator's decision stands between the transaction and commit.

Everything hinges on the prepared state. A prepared participant has given up its right to decide. It must survive a crash still holding the transaction, still holding its locks, and wait to be told. Whether real databases can do that is easy to check, because Postgres exposes the participant half of 2PC as ordinary SQL.
3.2Prepared transactions in Postgres
Postgres lets us play the participant's part ourselves. You run a transaction as normal, then end it with PREPARE TRANSACTION 'gid' instead of COMMIT, where the gid is a name you choose for the transaction (a global transaction identifier). Later, from any session (one connection to the database), COMMIT PREPARED 'gid' or ROLLBACK PREPARED 'gid' finishes it.
It's off by default: max_prepared_transactions is 0, and the documentation recommends leaving it there unless an external transaction manager, a program that plays the coordinator, is tracking every prepared transaction.
?How does a prepared transaction survive a crash?
Postgres records every change in its write-ahead log, or WAL, the append-only log from section 1.2, and flushes it to disk before a commit returns. A prepared transaction's state goes into the WAL too. Every Postgres transaction also gets a number, its XID, and Postgres keeps a dummy process entry for each prepared transaction so that its XID still counts as running. The header comment of the source file says it directly:
* Each global transaction is associated with a global transaction
* identifier (GID). The client assigns a GID to a postgres
* transaction with the PREPARE TRANSACTION command.
*
* We keep all active global transactions in a shared memory array.
* When the PREPARE TRANSACTION command is issued, the GID is
* reserved for the transaction in the array. This is done before
* a WAL entry is made, because the reservation checks for duplicate
* GIDs and aborts the transaction if there already is a global
* transaction in prepared state with the same GID.
*
* A global transaction (gxact) also has dummy PGPROC; this is what keeps
* the XID considered running by TransactionIdIsInProgress. It is also
* convenient as a PGPROC to hook the gxact's locks to.
*
* Information to recover prepared transactions in case of crash is
* now stored in WAL for the common case. In some cases there will be
* an extended period between preparing a GXACT and commit/abort, in
* which case we need to separately record prepared transaction data
* in permanent storage. This includes locking information, pending
* notifications etc. All that state information is written to the
* per-transaction state file in the pg_twophase directory."Keeps the XID considered running" is the sentence to remember. As far as the rest of the database is concerned, a prepared transaction is a transaction that hasn't finished and has no client attached. We'll see in section 3.3 what that costs.
The two durable steps show up in the WAL. The next experiment runs one 2PC transaction and then one ordinary transaction, each adding 1 to a row of a small test table called bench, and pg_waldump, a tool that prints a WAL's records, shows what each one wrote. In its output, rmgr names the part of Postgres that wrote the record, len (rec/tot) is the record's size in bytes, tx is the XID, and desc says what the record is.
begin; update bench set v = v + 1 where id = 7; prepare transaction 'demo-1';
commit prepared 'demo-1';
begin; update bench set v = v + 1 where id = 8; commit;rmgr: Heap len (rec/tot): 71/ 71, tx: 135415, desc: HOT_UPDATE ...
rmgr: Transaction len (rec/tot): 293/ 293, tx: 135415, desc: PREPARE gid demo-1: 2026-09-27 04:06:29.058304 UTC
rmgr: Transaction len (rec/tot): 42/ 42, tx: 0, desc: COMMIT_PREPARED 135415: 2026-09-27 04:06:29.061393 UTC
rmgr: Heap len (rec/tot): 71/ 71, tx: 135416, desc: HOT_UPDATE ...
rmgr: Transaction len (rec/tot): 34/ 34, tx: 135416, desc: COMMIT 2026-09-27 04:06:29.061850 UTCThe lsn and prev columns and the tail of each heap record are trimmed from this output. The Heap lines are the row updates themselves, one per transaction. The plain transaction (XID 135416) ends in one 34-byte COMMIT record. The 2PC transaction (XID 135415) needs a 293-byte PREPARE record, which holds the locks and state needed to rebuild it after a crash, and then a separate 42-byte COMMIT_PREPARED. (Its tx reads 0 because that record is written by whichever session runs COMMIT PREPARED, which doesn't need an XID of its own; the description names the XID it finishes, 135415.) Each of those records is flushed to disk before its command returns, so one commit has become two flushes, a cost we'll price in section 3.4.
3.3The blocking problem
Now the failure that 2PC is known for. The coordinator has collected yes votes from Orders and Payments, and then it crashes before it tells either of them the outcome.
Payments has prepared a transaction that debits an account. The coordinator crashes before sending COMMIT or ROLLBACK, and Payments restarts. What happens to the prepared transaction and its row lock?
Here is that run for order 2, drawn as state. Watch what each database is stuck holding.
We can watch a real Postgres do this. The next experiment plays Payments' part with a plain accounts table, where the work inside the transaction is subtracting 30 from account 1. It prepares the transaction, kills the server with pg_ctl stop -m immediate (an immediate stop with no clean shutdown, like a power cut) and starts it again. Then it looks at pg_prepared_xacts, the list of prepared transactions, and from another session tries to update the same row. lock_timeout makes a statement that has to wait for a lock give up after two seconds, where it would otherwise wait for ever. Finally pg_locks lists the locks that are currently held.
-- participant A, as the coordinator's connection
begin;
update accounts set balance = balance - 30 where id = 1;
prepare transaction 'xfer-42';
-- coordinator dies here. pg_ctl stop -m immediate, then start.
select gid, prepared, owner from pg_prepared_xacts;
-- any other session
set lock_timeout = '2s';
update accounts set balance = balance + 1 where id = 1;
select locktype, mode, granted, virtualtransaction, pid from pg_locks
where relation = 'accounts'::regclass or locktype = 'transactionid'; gid | prepared | owner
---------+------------------------------+----------
xfer-42 | 2026-09-27 04:04:12.72303+00 | postgres
(1 row)
Time: 2001.676 ms (00:02.002)
ERROR: canceling statement due to lock timeout
CONTEXT: while updating tuple (0,1) in relation "accounts"
locktype | mode | granted | virtualtransaction | pid
---------------+------------------+---------+--------------------+-----
transactionid | ExclusiveLock | t | -1/731 |
relation | RowExclusiveLock | t | -1/731 |The first table shows xfer-42 still prepared after the crash and restart. The other session's update waited its full 2,001 ms and then gave up with a lock timeout, which means the row is still locked by a transaction nobody is running. The last table shows who holds the locks, and it's the dummy entry from section 3.2: its virtualtransaction is -1/731 and its pid is empty. There's no backend (the server process Postgres runs for each connection) to cancel and no connection to kill, so pg_terminate_backend can't help you here.
Those locks are the visible half of the damage. The hidden half comes from VACUUM, the Postgres job that cleans up dead row versions, the old copies left behind whenever a row is updated or deleted. VACUUM can only remove a dead version if no running transaction could still need to see it, and a prepared transaction's XID counts as running. In the same test, updating 100,000 rows of an unrelated table and then vacuuming it printed the lines below, where a tuple is Postgres's word for one stored version of a row:
tuples: 0 removed, 200000 remain, 100000 are dead but not yet removable
removable cutoff: 731, which was 5 XIDs old when operation endedCutoff 731 is the prepared transaction's XID. After COMMIT PREPARED 'xfer-42', the same VACUUM reported 100000 removed. Left for weeks, a forgotten prepared transaction bloats every table in the database and, as the docs warn, can eventually force a shutdown to prevent transaction ID wraparound, which is what happens when Postgres's transaction counter gets too far ahead of its oldest unfinished transaction.
?Why does the coordinator's log make "no decision" mean abort?
Because commit happens at exactly one point: the coordinator durably writing its decision. If recovery finds no commit record for a transaction, the coordinator can't have told anyone to commit, so aborting is always consistent. This convention is called presumed abort, from IBM's R* system (Mohan, Lindsay and Obermarck, 1986), and it's what lets the coordinator skip the durable log write for aborts. The last frame of the diagram above is this rule at work.

Blocking is what happens when the coordinator fails. Even when nothing fails, 2PC isn't free.
3.4What 2PC costs
Every step of the protocol that has to survive a crash is a disk flush. A plain commit needs one, for its commit record. In 2PC every participant flushes a prepare record, then the coordinator flushes its decision, then every participant flushes a commit record, one after another with a network round trip between the phases. Through all of it, the participants hold their locks.
Here's the flush count on one Postgres, with the coordinator and the participant in the same place so that there's no network at all. The test commits a one-row update with a plain COMMIT, and then the same update with PREPARE TRANSACTION followed by COMMIT PREPARED, using pgbench (Postgres's benchmark tool) with one client for 10 seconds, three runs each:
| Commit path | Median latency | Runs | WAL flushes |
|---|---|---|---|
| BEGIN · UPDATE · COMMIT | 0.32 ms | 0.354 · 0.285 · 0.316 ms | 1 |
| BEGIN · UPDATE · PREPARE · COMMIT PREPARED | 0.79 ms | 0.792 · 0.721 · 0.988 ms | 2 |
That's roughly 2.5 times the commit latency on one node, before any network. The virtual disk in this test flushed in about 0.37 ms according to pg_test_fsync, far faster than a cloud volume, so on real storage every extra flush costs more. Across real machines, add roughly one network round trip per phase, and add the time the locks are held while the coordinator waits for the slowest participant.
Google's Spanner, a globally distributed database, gives the one large-scale public measurement. Its paper reports 2PC running over groups of replicated machines spread across three zones, 25 servers in each (a zone is Spanner's unit of deployment, a set of servers in one data centre), with a mean commit latency of 17.0 ms with one participant, 24.5 ms with two, 42.7 ms with 50, and 150.5 ms with 200. Those are means over ten runs. The authors call scaling to 50 participants "reasonable" and note that latency starts to rise noticeably at 100.
Those latencies matter because of the locks. If every commit touches the same row and holds its lock for the whole commit, a row can only commit about once per commit latency:
| 1 participant, 17.0 ms per commit | 1,000 ms ÷ 17.0 ms | ≈ 59 per second |
| 50 participants, 42.7 ms | 1,000 ms ÷ 42.7 ms | ≈ 23 per second |
| 200 participants, 150.5 ms | 1,000 ms ÷ 150.5 ms | ≈ 7 per second |
| most commits a single hot row can take, as participants are added | 59 → 7 / s | |
These figures are arithmetic on one simplifying assumption, and real locks are often held for longer than the commit itself, so treat them as a ceiling. They still show why the transaction that hurts is one that touches many participants and a popular row. All of this is what 2PC costs when nothing fails. The harder question is whether we can keep its atomicity and get rid of the blocking from section 3.3.
3.5Making 2PC non-blocking
2PC blocks because the decision lives in one place, the coordinator, and that place can fail. Before we look at the fixes we need one more idea. A consensus group is a small set of machines, usually three or five, that run a protocol such as Paxos or Raft to agree on one ordered log, so that the log survives the loss of any minority of them. Chapter 27 builds one. With that in hand, there are three ideas on the table.
| Approach | Idea | Used by |
|---|---|---|
| Three-phase commit (Skeen, 1981) | Add a pre-commit round so participants can finish without the coordinator | Probably nobody in production: it assumes a synchronous network, and a partition breaks it |
| Replicate the coordinator | Make the coordinator's log a consensus group, so its decision survives any one machine | Spanner, CockroachDB, YugabyteDB, TiDB |
| Paxos Commit (Gray and Lamport, 2006) | Store each participant's vote in its own consensus group; 2PC is the special case where each vote is kept on one machine only | The theory behind the second row |
?Why does replicating the coordinator fix blocking?
Because now "the coordinator crashed" means "the leader of a Paxos or Raft group crashed", and the group elects a new leader that can read the decision from the replicated log. Spanner is explicit about it: each participant is a Paxos group, and "running two-phase commit over Paxos mitigates the availability problems". The latencies in section 3.4 are what that costs.
Spanner's authors also make the argument for offering transactions at all: "We believe it is better to have application programmers deal with performance problems due to overuse of transactions as bottlenecks arise", instead of making everyone code around their absence.
3.6Where you'll meet 2PC
The table below needs a few names. XA is the standard interface through which a transaction manager drives 2PC across databases and message brokers that support it. JTA is Java's API for it, and Atomikos and Narayana are transaction managers that implement it. A shard is one of the pieces a large database splits its data into, each usually on its own machines.
| System | What it gives you | Watch out for |
|---|---|---|
XA (JTA, Atomikos, Narayana, MySQL XA START) | 2PC across databases and message brokers, driven by a transaction manager | The transaction manager's log is now a critical piece of state. If it's lost, participants may resolve transactions on their own guess, which XA calls a heuristic outcome |
Postgres PREPARE TRANSACTION | The participant half, for a manager you run | Orphans that hold locks and block VACUUM |
| Spanner, CockroachDB | 2PC inside the database, across its own shards, over consensus | Latency grows with the number of shards a transaction touches |
DynamoDB TransactWriteItems | Up to 100 items across tables, atomically (API reference) | One region, one service; not across DynamoDB and anything else |
| Kafka transactions | Atomic writes to many partitions plus consumer offsets | Only covers Kafka (section 7) |
Notice the pattern in that table: 2PC works well inside one system, where the same team owns the coordinator, the participants and the recovery code. It works badly between systems owned by different teams, where nobody owns the recovery.
That advice covers the case where you own both sides. For order 2, one side is a payment provider, and it won't be a participant in anybody's transaction.
04Undoing instead of rolling back
4.1Sagas
Most real checkouts can't use 2PC. The payment provider is another company's service, a program with its own database and its own team, reached over the network, and it won't join your XA transaction. Even if it would, holding locks while waiting on a third party is a bad idea, as the numbers in section 3.4 show. So we give up atomicity and replace it with compensation: undoing a step's effect with a new step, the way a refund undoes a charge.
A saga is a sequence of local transactions, each committed on its own. If a later step fails, we run compensating transactions that undo the earlier steps in business terms, in reverse order. The idea comes from Hector Garcia-Molina and Kenneth Salem's Sagas (SIGMOD 1987), written for long-lived transactions inside one database, and it maps directly onto services. Here's order 2 as a saga in which shipping fails. Each service has its own place, and a log underneath records which steps are done.
The log matters. A crashed orchestrator can read it on restart and resume where it stopped, rather than starting over.
?Why is a compensation not the same as a rollback?
A rollback makes it as if the transaction never happened. A compensation is a new transaction that happens after the original, and everyone may have seen the original in between. An email already went out, and a refund takes days to reach the card. Some actions can't be compensated at all: you can't unsend a notification or unship a parcel.
That's why the order of steps matters. Put the steps that can't be undone last. Chris Richardson's Microservices Patterns (Manning, 2018) calls the step after which the saga can only go forward the pivot transaction: before it, failures compensate, and after it, failures retry until they succeed.
4.2Orchestration or choreography
Someone has to run the saga. There are two shapes, and they differ in who knows what comes next. In the picture above, the saga log belongs to a single orchestrator, one piece of code that calls each step and decides what happens next. The alternative is choreography, where there's no central code: each service listens for an event, a message announcing that something happened such as OrderPlaced, does its own step, and publishes the next event.
| Orchestration | Choreography | |
|---|---|---|
| Who drives | One orchestrator calls each step and decides what's next | Each service reacts to the previous service's event |
| Where the state lives | The orchestrator's saga log | Spread across every service's events |
| Good for | Many steps, branching, compensations that must run in order | Two or three steps with a simple chain |
| Hard part | The orchestrator is a service you must make durable | Nobody can answer "what state is order 2 in?" |
| Tools | Temporal, AWS Step Functions, Camunda, a table and a worker | Kafka topics plus the outbox (section 5) |
Which should you pick? Teams usually end up with orchestration, because with choreography the saga's logic lives nowhere in particular. Adding a step means changing which events three services listen to, and debugging a stuck order means reading four services' logs. An orchestrator makes the flow a piece of code you can read, and its log answers "what state is it in" directly.
4.3What you give up: isolation
A saga is atomic in the end, since every step completes or is compensated. It isn't isolated, which means other readers can see its in-between states. Between T1 and C1 in the picture, the rest of the world could see a PENDING order for something that was about to be cancelled. These are the same kinds of anomaly that database isolation levels exist to prevent (chapter 19), and a saga has to handle them in application code:
| Anomaly | Example | Countermeasure |
|---|---|---|
| Dirty read | A report counts the PENDING order that's about to be cancelled | Semantic lock: a status field that readers respect |
| Lost update | C2 releases stock using a count read before another saga changed it | Commutative updates: stock = stock + 1, never "set to the old value" |
| Unrepeatable read | Two steps read the customer's credit limit and get different answers | Re-read and re-check in the step that commits |
| Compensation of a changed thing | C3 refunds a charge that support already refunded by hand | Make compensations idempotent, and check current state first |
Look at what a saga asks of every step. A step has to commit its own change and then get the next step started, which is a change to a database plus a message to someone, and that's the dual write from section 2 all over again. And the orchestrator will crash between sending a request and recording the reply, so on restart it sends the request again, which is the ambiguous timeout of section 2.2. We'll fix the first problem first.
05The transactional outbox
5.1Write the message as a row
We need a database row and a message to happen together, and 2PC between Postgres and Kafka isn't on offer. The way out is to turn the two writes into one. Instead of sending the event, the checkout inserts it into an outbox table, in the same local transaction as the business change:
BEGIN;
INSERT INTO orders (id, customer_id, total) VALUES (2, 7, 4999);
INSERT INTO outbox (aggregate_id, topic, payload)
VALUES (2, 'orders', '{"type":"OrderPlaced","order_id":2}');
COMMIT;Now there's one commit record in one log. Either the order and its event both exist, or neither does. A separate process, the relay, reads the outbox table and publishes each row to Kafka. Chris Richardson writes the pattern up on microservices.io. Here's order 2 going through it, including the failure that gives the pattern its limit:
?Why doesn't the outbox give exactly-once delivery?
Because the relay has its own dual write: publish to Kafka, then mark the row sent. It's the same problem, moved to a place where it does little harm. The relay always publishes before marking, so a crash produces a duplicate and never a loss, and a duplicate is something a receiver can detect.
5.2The consumer's half: the inbox
Detecting the duplicate is the job of the consumer, the service that reads the messages, here the warehouse. A robust design mirrors the outbox. The consumer records each event's ID in a table in its own database, in the same local transaction as the effect of handling the event. This table is called the inbox.
BEGIN;
INSERT INTO processed_events (event_id) VALUES ($1)
ON CONFLICT DO NOTHING
RETURNING event_id; -- no row back? we've done this one; stop
-- only if the insert returned a row:
UPDATE stock SET reserved = reserved + 1 WHERE sku = 'A-17';
COMMIT;When the second copy of the OrderPlaced event arrives, the insert finds the ID already there, does nothing and returns no row, so the application skips the update. Because the ID and the stock change commit together, a crash between them is impossible: either both happened or neither did. This is the same trick Kafka's own docs recommend for external systems (section 7.3): store "what I've processed" in the same place, and the same transaction, as the output.
5.3Polling relay or change data capture
How does the relay find unsent rows? There are two common ways, and both leave a different thing to watch. Change data capture, or CDC, means reading the database's own log of changes (its WAL in Postgres) instead of querying its tables. Postgres offers this through logical decoding, which turns the WAL into a stream of row-level changes that a program can follow. A program reading the stream holds a replication slot, Postgres's bookmark recording how far that reader has got. MySQL's equivalent of the WAL is called the binlog.
| Polling relay | Change data capture (CDC) | |
|---|---|---|
| How | SELECT ... FOR UPDATE SKIP LOCKED on unsent rows, publish, mark sent | Read the database's WAL (Postgres logical decoding, MySQL binlog) and emit outbox inserts |
| Tool | A loop you write | Debezium's outbox event router |
| Latency | Your polling interval | Close to commit time |
| Load on the database | A query per poll, plus index churn on the outbox | Reads the WAL it already writes; holds a replication slot |
| Failure to watch | Outbox growing because the relay died | A stalled slot retaining WAL until the disk fills |

A polling relay that several workers can run at once looks like this:
BEGIN;
WITH batch AS (
SELECT id FROM outbox
WHERE sent_at IS NULL
ORDER BY id
LIMIT 100
FOR UPDATE SKIP LOCKED -- other relays skip rows this one holds
)
SELECT o.* FROM outbox o JOIN batch USING (id);
-- publish every row, wait for acks, then:
UPDATE outbox SET sent_at = now() WHERE id = ANY($1);
COMMIT;The BEGIN matters, because the row locks that FOR UPDATE takes last only until the transaction ends.
?Why SKIP LOCKED?
Without it, a second relay's FOR UPDATE waits on the first relay's rows and the relays run one at a time. With it, each relay takes whatever rows nobody else holds. The cost is ordering: two relays can publish neighbouring events out of order. If per-order ordering matters, give each relay a partition of aggregate_id instead.
With CDC the thing to remember is the replication slot. The WAL is stored as a series of files, and Postgres can't delete a file until every slot has confirmed reading past it. So if the CDC connector (the program reading the slot) stops, the WAL files pile up and the database server's disk fills. Alert on pg_replication_slots lag, and set max_slot_wal_keep_size so that a dead connector costs you the slot instead of the database.
The outbox and inbox keep messages between your own services from being lost or doubled. Requests that arrive from outside, from a browser or a partner's server, have the same problem at the front door.
06Idempotency keys
6.1How the server uses the key
When a client's call to your API times out, it retries, and from section 2.2 you know it can't tell whether the first attempt landed. The answer is to make repeating the call harmless. An operation is idempotent if doing it twice has the same effect as doing it once. A charge isn't idempotent by nature, so we make it so with an idempotency key: a unique value the client generates once per logical operation and sends with every retry of it. The server remembers which keys it has seen. The inbox of section 5.2 is the same idea applied to messages.
Stripe's is the design most people copy. Its API reference says it works "by saving the resulting status code and body of the first request made for any given idempotency key, regardless of whether it succeeds or fails. Subsequent requests with the same key return the same result, including 500 errors."
Four decisions go into a design like that:
| Decision | Stripe's answer | Why |
|---|---|---|
| Who makes the key | The client, ideally a V4 UUID (a random 128-bit ID), up to 255 characters | Only the client knows two requests are the same operation |
| Same key, different body | Error | A reused key with new parameters is a bug, not a retry |
| How long it's kept | Keys can be pruned once they're at least 24 hours old | Retries happen in seconds or minutes; storage isn't free |
| What's stored | The status and body, once the endpoint starts executing | A retry must get the same answer, not a fresh one |
DynamoDB makes different choices for the same problem. A ClientRequestToken on TransactWriteItems is valid for 10 minutes after the first request completes, and a reuse with different parameters returns IdempotentParameterMismatch (API reference).
?Why generate the key on the client and not the server?
Because the server can't tell a retry from a second purchase. Two identical POST /charges for $49.99 might be one customer clicking twice, or one request retried after a timeout. The client knows which it meant, so the client names the operation, generates the key before the first attempt, and reuses it for every retry of that attempt.
6.2Two retries at once
Harder than a retry after the first request finished is a retry that arrives while the first is still running, because the client's timeout fired early. Both requests look up the key and find nothing, and both go ahead and charge. So we make "claim this key" an atomic insert, which lets the database decide the race.
The next experiment sends two requests with the same key, started about 200 ms apart. The table uses the key as its primary key, so Postgres refuses to store a second row with the same key. on conflict (key) do nothing turns that refusal into "insert nothing", and returning 'inserted' prints a row only when an insert happened. pg_sleep(1) stands in for a second of real work, such as calling the card network. In the output each line is prefixed with the request that printed it, A or B, and the two timestamps per request are start and after_insert. The blank rows that pg_sleep prints have been removed.
create table idempotency_keys (
key text primary key,
request_hash text not null,
status int,
response jsonb,
created_at timestamptz not null default now()
);
-- each request runs:
begin;
select clock_timestamp()::time(3) as start;
insert into idempotency_keys (key, request_hash)
values ('k-7f3a', 'sha256:amount=500')
on conflict (key) do nothing
returning 'inserted';
select clock_timestamp()::time(3) as after_insert;
select pg_sleep(1); -- stands in for the real work
commit;A: BEGIN
A: 04:06:39.959
A: inserted
A: INSERT 0 1
A: 04:06:39.962
A: COMMIT
B: BEGIN
B: 04:06:40.157
B: INSERT 0 0
B: 04:06:40.967
B: COMMITA's insert took 3 ms (39.959 to 39.962) and printed inserted, so A owns the key and goes on to do the work. B started 198 ms later, and its insert waited roughly 810 ms (40.157 to 40.967) on A's uncommitted row in the primary key index, until A committed. Then it inserted nothing (INSERT 0 0). B knows it doesn't own the key, so it doesn't charge the card. It returns A's saved response, or a 409 telling the client to retry shortly if A hasn't finished.
6.3When the work calls someone else
That works when all the work is one local transaction. It gets harder when the handler itself calls a third party: create the charge row, call the card network, send a receipt email. Why not hold one database transaction across the whole request? Because the external call can take seconds, and the transaction would hold its locks the whole time. Worse, if the call succeeds and your commit fails, you've charged the card with no record of it.
Brandur Leach's Implementing Stripe-like Idempotency Keys in Postgres solves it with recovery points. Each key's row records how far the request got. In his ride-sharing example the points are started, ride_created, charge_created and finished. Each local step commits and advances the recovery point atomically, and a retry resumes from the last recorded point. Committing a recovery point before each external call means there's always a durable note saying "this might have happened; check before redoing it". Each foreign call must itself be idempotent, which usually means passing your own idempotency key through to the provider. It's a small saga, with the key's row as its log.
We've now built every piece: outbox and inbox for messages, keys for requests. They all rest on the same bargain of at-least-once delivery plus duplicate detection. That bargain is what messaging systems advertise as "exactly-once", and it deserves a careful reading.
07What "exactly-once" can mean
7.1Delivery versus processing
A message crossing a network can be delivered at most once, where the sender never retries and so may lose the message, or at least once, where it retries until acknowledged and so may duplicate it. The ambiguous timeout from section 2.2 means there's no third option at the delivery layer.
What systems can offer is exactly-once effect: every message is delivered at least once, and duplicates are detected and discarded before they change anything. That's what the outbox with an inbox, and idempotency keys, already built. The same three choices show up in how a consumer commits its position in a stream. Kafka splits each topic into partitions, each an append-only log of messages in which a message's position is its offset. A consumer reads messages in order and remembers how far it has got by committing an offset.
| Guarantee | Mechanism | What can still go wrong |
|---|---|---|
| At most once | Commit the offset, then process | A crash after commit loses the message |
| At least once | Process, then commit the offset | A crash before commit reprocesses it |
| Exactly-once effect | At least once, plus dedupe stored atomically with the effect | Any side effect outside that atomic store, like an email or an HTTP call |
7.2What Kafka's exactly-once is
Kafka 0.11 (2017, KIP-98) added two features, and since Kafka 3.0 the first is on by default. The producer is the program that writes messages to Kafka, and the class's own header comment describes both:
* From Kafka 0.11, the KafkaProducer supports two additional modes: the idempotent producer and the transactional producer.
* The idempotent producer strengthens Kafka's delivery semantics from at least once to exactly once delivery. In particular
* producer retries will no longer introduce duplicates. The transactional producer allows an application to send messages
* to multiple partitions (and topics!) atomically.
* </p>
* <p>
* From Kafka 3.0, the <code>enable.idempotence</code> configuration defaults to true. When enabling idempotence,
* <code>retries</code> config will default to <code>Integer.MAX_VALUE</code> and the <code>acks</code> config will
* default to <code>all</code>.
* ...
* To take advantage of the idempotent producer, it is imperative to avoid application level re-sends since these cannot
* be de-duplicated.
* ...
* Finally, the producer can only guarantee idempotence for messages sent within a single session.The idempotent producer is an idempotency key built into the protocol. Each producer gets an ID from the broker (a Kafka server that stores partitions), the producer numbers its batches per partition, and the broker discards a batch whose sequence number it has already written. It fixes duplicates from the producer's own retries, within one producer session, and nothing else.
The transactional producer is 2PC with Kafka as every participant. The coordinator role from section 3 is played by one of the brokers, called the transaction coordinator, which keeps its decisions in an internal Kafka topic. Transactions are what make consume-transform-produce atomic: read a batch of input, process it, write the output. One more problem has to be handled first. A transactional producer identifies itself with a fixed name, its transactional.id, so that a replacement can pick up where a crashed instance left off. But if a producer only stalls and a replacement starts under the same name, the old instance may wake up and keep writing, a zombie. So each time a producer starts under that name, Kafka gives it a higher epoch number. The brokers reject any write carrying an older epoch, and shutting out the stale instance this way is called fencing.
transactional.id, the coordinator bumps the producer's epoch. Any older instance with the same ID is now fenced and its writes will be rejected.Why do the offsets have to be in the transaction? Because "I've processed input up to offset 500" and "here's the output for it" are the two writes of a dual write. Kafka can make them atomic only because both are writes to Kafka. If the transaction aborts, the output is invisible to read_committed readers and the offset reverts, so the batch is processed again and produces the same output.
7.3Where it stops
Once your processing touches anything outside Kafka (a database write, an HTTP call, an email), those effects aren't in the transaction. Kafka's own design documentation says so directly:
"When writing to an external system, the limitation is in the need to coordinate the consumer's position with what is ... stored as output. The classic way of achieving this would be to introduce a two-phase commit between the storage of the consumer position and the storage of the consumers output. But this can be handled ... by letting the consumer store its offset in the same place as its output."
That's the inbox from section 5.2, recommended by Kafka itself. The same page goes on: "Exactly-once delivery for other destination systems generally requires cooperation with such systems ... Otherwise, Kafka guarantees at-least-once delivery by default."
Pat Helland's Life beyond Distributed Transactions (CIDR 2007) is the long version of this chapter's argument. At scale, he writes, applications must assume messages are retried and reordered, keep each atomic unit to one "entity", and make every message handler idempotent. Most of the patterns in this chapter are in that paper.
08Choosing and operating
8.1Which tool for which problem
Each tool answers one of the failures from the checkout. Here they are as a decision table:
| Situation | Reach for | Not |
|---|---|---|
| Both sides can live in one database | A local transaction | Any of the rest |
| Shards of one distributed database | Its built-in transactions (Spanner, CockroachDB, DynamoDB) | Hand-rolled 2PC |
| A database write plus a message | Outbox, with an inbox on the consumer | A dual write, or XA with the broker |
| A multi-step business process across services | An orchestrated saga | Choreography past three steps |
| A client calling your API that changes state | Idempotency keys, stored with the effect | "We'll dedupe in the logs" |
| Kafka in, Kafka out | Kafka transactions, read_committed downstream | Committing offsets by hand |
| Kafka in, database out | Offsets or event IDs stored in the database | Kafka transactions alone |
8.2What to watch
Every one of these patterns leaves a table that grows when something is stuck. Monitor those tables, not just error rates:
-- Is any 2PC participant orphaned? (section 3.3) Anything prepared for more than a minute.
SELECT gid, prepared, now() - prepared AS age
FROM pg_prepared_xacts WHERE prepared < now() - interval '1 minute';
-- Are sagas stuck mid-flight? (section 4)
SELECT state, count(*) FROM sagas
WHERE updated_at < now() - interval '15 minutes' AND state NOT IN ('DONE','COMPENSATED')
GROUP BY state;
-- Is the outbox draining? (section 5.1) The oldest unsent event.
SELECT now() - min(created_at) AS oldest_unsent FROM outbox WHERE sent_at IS NULL;
-- Is the CDC relay holding back WAL? (section 5.3)
SELECT slot_name, active,
pg_size_pretty(pg_wal_lsn_diff(pg_current_wal_lsn(), restart_lsn)) AS retained
FROM pg_replication_slots;8.3Rules that hold up
- Keep both sides of an invariant in one database when you can. A local transaction beats every pattern in this chapter.
- Never write to two systems in one request path. Write one row, and let a relay deliver the rest.
- Assume every message arrives at least once. Record the event ID in the same transaction as the effect.
- Generate the idempotency key on the client, and store it next to the effect. A key in a different system is a dual write.
- Order saga steps so the ones that can't be undone come last, and write compensations as relative changes that are safe to repeat.
- Resolve an orphaned prepared transaction from the coordinator's log. Don't guess.
8.4What you trade for what
| You get | You pay | When the bill arrives |
|---|---|---|
| 2PC's atomic commit across databases | Extra flushes, locks held across round trips, blocking if the coordinator dies | As latency on hot rows, or an orphaned prepared transaction |
| A saga's availability, with no locks across services | No isolation, and compensations that can't undo everything | As reports that show orders that were later cancelled |
| The outbox's guarantee that an event follows its row | At-least-once delivery and a relay to run | As duplicates, an outbox that grows, or a stalled replication slot |
| An idempotency key's safe retries | A key table in the same database, and a retention policy | As a double charge when the key lives somewhere else |
| Kafka's exactly-once between topics | It stops at Kafka's edge | As duplicate side effects in whatever the consumer writes to |
8.5Symptom, cause, fix
| Symptom | Likely cause | Fix |
|---|---|---|
Rows locked with no session holding them; VACUUM removes nothing | Orphaned prepared transaction | Apply the coordinator's recorded decision with COMMIT PREPARED or ROLLBACK PREPARED; alert on pg_prepared_xacts age |
| Order exists but downstream never heard about it | Dual write; process died between the two | Outbox |
| Duplicate charges after a provider timeout | Retry without an idempotency key | Client-generated key, stored atomically with the charge |
| Duplicate side effects after a consumer restart | At-least-once delivery with no dedupe | Inbox table or offsets stored with the output |
| Database server's disk filling with WAL files | Stalled CDC replication slot | Restart the connector; max_slot_wal_keep_size |
Order stuck in PENDING for hours | Saga orchestrator crashed between steps, or a step is failing forever | Durable saga log with resume; alert on age; dead-letter after N retries |
| Reports briefly show orders that were then cancelled | Sagas aren't isolated | Semantic lock: filter on status, or report from committed end states |
09Summary
- A dual write has no atomic outcome. Each system commits to its own log, and the process can die between them, which is how order 2 was charged and never saved.
- A timeout doesn't tell you whether it happened. Every technique here makes that ambiguity safe: find out, undo, or repeat harmlessly.
- 2PC commits at one record: the coordinator's decision. Participants promise first, and that promise must survive a crash.
- A prepared participant blocks until it hears the decision. In Postgres it keeps its locks with no backend attached and stops
VACUUMat its XID. - 2PC costs extra flushes and held locks. About 2.5 times the commit latency on one node, and 17.0 to 42.7 ms at 1 to 50 participants in Spanner.
- Replicating the coordinator removes blocking. That's why Spanner and CockroachDB run 2PC over consensus, and why 2PC works best inside one system.
- A saga trades isolation for availability. Compensations add new history, must be idempotent, and can't undo everything.
- The outbox turns two writes into one. The relay delivers at least once, so consumers deduplicate in the same transaction as their effect.
- Idempotency keys are claimed with an atomic insert. Store them next to the effect they protect, not in a separate system.
- "Exactly-once" means at-least-once plus deduplication. Kafka's version covers Kafka to Kafka, and anything outside needs its own dedupe.
10Build this
A checkout that survives being killed.
- Start two Postgres instances on different ports,
ordersin one andpaymentsin the other. Write the naive dual write, then kill the process withkill -9at random points in a loop and count the mismatched rows. - Replace it with an outbox and a polling relay using
SKIP LOCKED, and a consumer with an inbox table. Run the same kill loop against the handler, the relay and the consumer. The mismatch count should be zero, and the consumer should log the duplicates it discarded. - Add an idempotency-key table to the payments side. Fire each request twice, 50 ms apart, from different connections, and check that exactly one charge row exists per key.
- Finally, set
max_prepared_transactions, do the transfer withPREPARE TRANSACTION, kill the coordinator between phases, and write the recovery script that reads your coordinator log and resolves every row inpg_prepared_xacts.
11Interview questions
beginnerWhat is the dual-write problem?›
Writing to two systems (say, a database and a message broker) as two separate operations. Each is atomic on its own, but the process can crash, or one call can time out, between them. You end up with a row and no event, or an event and no row. Usually the fix is the transactional outbox: write the event as a row in the same local transaction, and publish it afterwards.
beginnerWhat's an idempotency key, and who generates it?›
A unique value naming one logical operation, sent with every retry of it. Your server stores it with the result, atomically with the effect, and returns the stored result for any repeat. The client generates it, because only the client knows whether two identical requests are a retry or two separate purchases. Stripe keeps keys for at least 24 hours and rejects a reused key with different parameters.
intermediateWalk through two-phase commit. Where exactly does the transaction commit?›
Phase one: the coordinator sends PREPARE, and each participant makes its changes durable, keeps its locks and votes yes or no. A yes is a binding promise. If all vote yes, the coordinator durably logs COMMIT. That log record is the commit point. Phase two: it tells each participant, retrying until they all acknowledge. With presumed abort, a coordinator that finds no commit record on recovery aborts.
intermediateWhy is 2PC called a blocking protocol?›
A participant that has voted yes can't decide on its own. If the coordinator fails before telling it the outcome, it must keep the transaction prepared, with its locks, until the coordinator recovers. In Postgres that's a row in pg_prepared_xacts, locks with no backend PID, and VACUUM unable to clean anything newer than its XID. Systems like Spanner avoid this by making the coordinator a Paxos group.
intermediateHow is a saga different from a distributed transaction?›
A saga is a series of local transactions that each commit immediately, with compensating transactions to semantically undo earlier steps if a later one fails. It never holds locks across services, so it's available and fast, but it isn't isolated: other readers see intermediate states. Compensations are new transactions, not rollbacks, and some actions (an email, a shipment) can't be compensated, so they go after the pivot step.
deepKafka claims exactly-once semantics. What does it guarantee?›
Two things. The idempotent producer (default since 3.0) uses a producer ID and per-partition sequence numbers so its own retries don't duplicate, within one session. Transactions let a consume-transform-produce app write output and input offsets atomically, with read_committed consumers seeing only committed data, and fence zombie instances by epoch. It covers Kafka to Kafka. Writes to a database or calls to an API are at-least-once unless you store offsets or event IDs with the output.
deepTwo retries with the same idempotency key arrive at the same time. How do you stop both from charging?›
Claim the key with an atomic insert on a unique index, such as INSERT ... ON CONFLICT DO NOTHING RETURNING, in the same database as the charge. The second insert blocks on the first's uncommitted index entry, then inserts nothing, so exactly one request owns the key. Whoever loses returns the saved response, or a 409 asking the client to retry if the first is still running. A check-then-insert in application code, or a key stored in a different system, both race.
deepYou find a prepared transaction in Postgres that's three days old. What do you do?›
Don't guess. Find which coordinator created it (the GID usually encodes it) and look up its decision in that coordinator's log. If it logged commit, run COMMIT PREPARED. If there's no commit record, presumed abort says ROLLBACK PREPARED is safe. Then find out why recovery didn't resolve it automatically, add an alert on pg_prepared_xacts age, and check table bloat, since VACUUM couldn't remove anything newer than its XID for three days.
12Go deeper
Which single record is the commit point of a two-phase commit?›
The coordinator's durable COMMIT decision. Before it, any participant can still be aborted; after it, every participant must eventually commit.
A prepared transaction in Postgres shows no pid in pg_locks. Why?›
It has no backend. Postgres attaches its locks to a dummy PGPROC so the XID stays "in progress", and that's also why it survives restarts.
Your outbox relay publishes, then crashes before marking the row sent. What happens?›
It publishes the row again after restart. The outbox gives at-least-once delivery, so consumers must dedupe by event ID.
Does enable.idempotence stop your application from sending the same order twice?›
No. It dedupes the producer's own internal retries within one session. If your code calls send() twice, Kafka stores two messages.
Lost messages, timeouts and retries, and why NFS is built from idempotent operations so that a repeated request does no harm: the same problem as section 2.2, from the operating system side. Free online at ostep.org.
The original paper: long-lived transactions as sequences with compensations. Short, and still the clearest statement of what you give up. PDF.
Shows 2PC as a degenerate consensus protocol and builds Paxos Commit, the non-blocking version. Microsoft Research.
2PC over Paxos in production, with the participant-count latency table used in section 3.4. Google Research.
Entities, activities and idempotent messages: the design philosophy behind sagas and outboxes at scale. PDF.
How prepared transactions are stored, recovered and finished. The header comment quoted in section 3.2 is the map. REL_16_4.
Recovery points, atomic phases and a complete schema for keys that protect calls to third parties. brandur.org.
The Kafka design proposal: producer IDs, epochs, the transaction coordinator and control markers. Apache wiki.
The CDC implementation of the outbox pattern, with the table layout it expects. Docs.
13Related chapters
The replicated log that turns 2PC's coordinator from a single point of failure into a group. Chapter 27.
What a single-database transaction guarantees, and the isolation a saga gives up. Chapter 19.
Why MULTI/EXEC isn't a transaction, and why an idempotency key in Redis
doesn't protect a write in Postgres. Chapter 22.
The locks a prepared transaction keeps holding, and what waiting on them costs. Chapter 13.