A single business operation must write rows on two different shards atomically. Explain why teams avoid two-phase commit (2PC/XA) for this and what they build instead.
answer
- Prepare holds locks until the coordinator decides
- Coordinator dies → in-doubt, blocked, horizon pinned
- Availability = product of all participants
- First choice: make it single-shard
- Saga + compensation + outbox + idempotency key
basics
~20 sTwo-phase commit holds locks while it waits for a coordinator, blocks in-doubt if the coordinator dies, adds round trips and fsyncs, and makes availability the product of all participants. Teams instead redesign so the write is single-shard, or use a saga with idempotent compensating steps plus a transactional outbox.
solid answer
~60 s**Why 2PC hurts:** the participant shards `PREPARE`, durably log, and then hold their locks until the coordinator says commit or abort. If the coordinator dies in that window, the transactions are *in-doubt* — locks held, rows unreadable by conflicting readers, and nothing can resolve them without recovery. The commit costs extra round trips and durable writes on the critical path, and the operation only succeeds if every participant is up, so availability multiplies down. Operationally, orphaned prepared transactions are a well-known way to pin a database's oldest-transaction horizon and block cleanup. **What we do instead, in order:** 1. **Make it single-shard.** Choose the shard key or denormalize so both writes land on one shard and one ordinary local ACID transaction covers them. This is the real answer most of the time. 2. **Saga.** Sequence local transactions, each with a compensating action, driven by a durable state machine. 3. **Transactional outbox.** Write the business row and an outbox row in one local transaction; a relay publishes it at-least-once; the downstream shard applies it idempotently keyed by a message id. You trade atomicity for eventual consistency, so you need idempotency keys, retries, and a reconciliation job.
code
sql · 14 linesBEGIN;
UPDATE accounts
SET balance_cents = balance_cents - 5000
WHERE tenant_id = 42 AND account_id = 7
AND balance_cents >= 5000;
INSERT INTO outbox (message_id, topic, payload, created_at)
VALUES ('3f9c-...-a1', 'transfer.debited',
'{"transfer_id":"3f9c-...-a1","to_account":991,"amount_cents":5000}',
now());
COMMIT;
-- a relay reads outbox and publishes; delivery is at-least-oncego deeper
Know that a transaction cannot span two independent databases for free, that two-phase commit exists, and that the practical fix is usually to arrange the data so both writes land on one shard.
Explain prepare/commit, the in-doubt blocking window, and the latency and availability costs, then name saga plus outbox plus idempotency as the standard alternative.
Show the whole production shape: co-locate first, saga with compensations and a durable state machine, outbox with at-least-once delivery, idempotent apply keyed by message id, reconciliation, and alerts on stuck sagas and prepared-transaction age.
Treat this as a consistency-model decision with product consequences — which invariants are genuinely global, what intermediate states the domain must expose, and where the narrow remaining case for a coordinated commit sits.
## The requirement "Debit account A, credit account B" where A and B hash to different shards. Inside one database this is one transaction and it is over. Across two independent database servers there is no shared transaction log, so atomicity must be manufactured. ## How two-phase commit works A coordinator drives two rounds: 1. **Prepare.** Each participant executes its work, makes it durable, and replies *prepared* — a promise it can still commit even after a crash. Its locks stay held. 2. **Commit / abort.** If every participant voted yes, the coordinator durably records the decision and tells everyone to commit; any no means abort. The protocol is correct: it never lets one side commit while the other aborts. Its cost is where it hurts. ## Why it is avoided **It blocks.** If the coordinator crashes after prepare and before the decision reaches participants, those participants are *in-doubt*. They cannot unilaterally commit (the decision might have been abort) and cannot abort (it might have been commit). They sit holding locks. Any transaction that needs those rows waits. A prepared transaction left behind also pins the engine's oldest-active-transaction horizon, which stops old row versions and log segments from being reclaimed — a slow-motion outage that fills a disk hours later. This is why operators alarm on prepared-transaction age. **It is slow on the critical path.** Two network round trips plus a durable log flush at each participant plus one at the coordinator, all inside the user's request. On top of that, the lock hold time is the whole protocol duration rather than the work duration, so conflict rates rise and the system's effective concurrency drops precisely when it is busiest. **Availability multiplies down.** The transaction succeeds only if every participant *and* the coordinator are healthy for its duration. With three participants at 99.9% each, the compound path is well under the availability of any one of them. Sharding was supposed to make failures independent; 2PC re-couples them. **Operationally it is a second system.** The coordinator needs durable, replicated storage and its own failover, plus recovery tooling to resolve in-doubt transactions. Some engines' distributed-transaction support interacts awkwardly with replication and connection pooling — a pooled connection cannot be handed to another session while a transaction is prepared on it. None of this makes 2PC wrong in principle. It makes it a poor fit for a high-rate, user-facing write path. It remains reasonable for low-rate control-plane operations where correctness dominates and latency does not. ## What replaces it ### 1. Eliminate the cross-shard write The cheapest distributed transaction is the one you do not have. If both rows can share a shard key — same tenant, same user, same account group — the operation becomes one local ACID transaction. Deliberately co-locating the entities that mutate together is the single highest-leverage decision in a sharded schema. Where a global uniqueness or balance invariant forces two owners, consider whether the invariant can be pushed onto one shard (for example, a per-tenant ledger rather than a global one). ### 2. Saga with compensations Model the operation as a sequence of local transactions, each with a semantic **compensating** transaction that undoes its business effect (refund, release reservation, cancel). A durable orchestrator — itself a table of saga state — advances the steps and, on failure, runs the compensations in reverse. What sagas are not: rollback. Intermediate states are visible to other readers, so the design must tolerate seeing a debit before its credit. That is a product decision as much as an engineering one: "pending", "reserved", and "settling" states usually have to appear in the domain model. ### 3. Reservation / two-step pattern A business-level analogue of prepare: first a local transaction reserves capacity or funds with an expiry (a hold row), then a second confirms it. The expiry is the crucial difference from 2PC — an abandoned reservation cleans itself up instead of blocking forever, because it holds application state rather than database locks. ### 4. Transactional outbox plus idempotent apply Within shard A's local transaction, insert both the business row and a row into an `outbox` table. A relay process reads the outbox and publishes the message, then marks it sent. Delivery is at-least-once, so shard B must apply idempotently: carry a stable message or idempotency key, record applied keys in the same transaction as the effect, and treat a duplicate as a no-op. This eliminates the dual-write problem — the classic bug where you write the database, then the message broker, and crash in between. ### 5. Reconcile Eventual consistency needs an auditor. A periodic job compares the two sides — sum of ledger entries versus balances, orders versus reservations — reports drift, and repairs or alerts. Without it, small leaks accumulate invisibly. ## What you must have either way Idempotency keys on every externally triggered write, bounded retries with backoff, a durable record of in-flight operations, timeouts on every step, and monitoring for stuck sagas. Whether you keep 2PC for a narrow control-plane case or go all-saga, those are the price of admission.
- How is a saga different from a database rollback?A rollback erases uncommitted changes so no other session ever saw them. A saga's steps each commit locally and are immediately visible, so undoing one means running a new compensating transaction that reverses its business effect — a refund rather than an un-charge. That means intermediate states must be legal in the domain model, and compensations must themselves be idempotent and able to handle the case where the original step partially succeeded.
- What does an at-least-once outbox demand of the consuming shard?Idempotency. The relay can crash after publishing but before marking the row sent, so the same message will be delivered again. The consumer needs a stable message or idempotency key, must record that key in the same local transaction as the effect it applies, and must treat a repeat as a no-op. It should also tolerate out-of-order arrival, usually by carrying a version or sequence number and ignoring anything not newer than what it already applied.
- When is 2PC still a reasonable choice?Low-rate operations where correctness dominates and a few hundred milliseconds do not matter — control-plane changes, provisioning, occasional administrative moves — and where you have a properly replicated coordinator plus monitoring on prepared-transaction age. It is the high-QPS user-facing path that cannot afford the lock hold times, the compounded availability, and the blast radius of an in-doubt transaction.
2PC is a group of people each holding a door open until the organizer shouts go — if the organizer disappears, everyone stands there. A saga is each person walking through and, if the plan falls apart, walking back out on their own.
saying these in an interview costs you the question
- Claiming 2PC gives the same availability as a local transaction because 'it's still ACID'
- Not knowing that a prepared transaction holds locks and can block indefinitely if the coordinator is lost
- Describing a saga as a rollback, and assuming intermediate states are invisible to other readers
- Writing the database and then publishing to a broker as two separate steps, with no outbox and no idempotency key
- Skipping reconciliation on the grounds that retries make eventual consistency exact