A distributed transaction is any change whose parts live in different failure domains, such as two database shards, two services or a database and a message broker, and which you need to succeed or fail as a unit. The textbook answer is two-phase commit. The practical answer is a ladder of options, and the engineering is choosing the lowest rung that gives the guarantee you need.
This article explains what makes cross-node atomicity hard, how the classic protocol fails, how distributed databases hide most of the cost, and what to do across services that cannot share a transaction, ending with a worked money transfer and a checklist.
Why atomicity across machines is hard
On one machine, commit is a single durable write: a log record either reached disk or it did not, and recovery replays the log. Across machines there is no single write. Each participant can make its own part durable, but none of them knows whether the others did, and messages can be delayed or lost. The participants need to agree on one outcome, commit or abort, and a participant that has promised to commit cannot change its mind, because the others may already have acted on that promise.
That promise creates the central problem, the in-doubt window. Between voting yes and learning the outcome, a participant must keep the change ready to commit and keep the locks that protect it. If the node deciding the outcome crashes in that window, the participant can do nothing safely: committing might contradict an abort, and aborting might contradict a commit. Every technique below either makes that window short, makes it recoverable without a human, or avoids it by giving up isolation.
The decision ladder
Work down this table and stop at the first rung that fits. Each rung down buys independence and costs guarantees your code must then provide.
| Rung | When it fits | What you get | What you pay |
|---|---|---|---|
| 1. Do not distribute | Data that changes together can live in one database or one partition. | An ordinary local ACID transaction. | A coarser ownership boundary; one database to scale. |
| 2. Distributed SQL database | One logical database spread over many nodes (Spanner, CockroachDB, TiDB, YugabyteDB). | Serializable or snapshot transactions across shards; the database runs the commit protocol. | Cross-region latency on commit; retries on contention. |
| 3. XA / 2PC across resources | Two different databases or a database and a queue, operated by one team, low volume. | Atomic commit across heterogeneous systems. | Blocking on coordinator failure, in-doubt transactions, lock time across round trips. |
| 4. Saga | Steps owned by different services, long-running, each reversible. | Eventual all-or-compensated outcome. | No isolation: other readers see intermediate states. See the saga pattern. |
| 5. Outbox plus idempotent consumer | One local change must reliably cause another service's change. | At-least-once delivery with exactly-once effect. | Asynchrony; the second change lands later. See transactional outbox. |
Rung 1 is underrated. Many 'distributed transaction' problems disappear when the order, its lines and its payment intent share a partition key, or when a counter moves to the service that owns the invariant.
Two-phase commit and XA, as they behave in production
In two-phase commit a coordinator asks every participant to prepare. A participant that votes yes has written its changes and its vote durably and still holds its locks. If every vote is yes, the coordinator durably records commit and tells everyone; any no vote or timeout before the decision leads to abort. The rule that keeps it correct is that the coordinator logs its decision before sending it, so recovery can finish the job. XA is the standard interface that lets a transaction manager drive this protocol across databases and brokers.
PostgreSQL exposes the participant side directly, which makes the mechanics easy to see:
-- postgresql.conf on every participant: max_prepared_transactions = 64 (default 0 = disabled)
-- Participant 1 (orders DB), driven by the coordinator
BEGIN;
UPDATE stock SET reserved = reserved + 1 WHERE sku = 'X-1' AND reserved < on_hand;
PREPARE TRANSACTION 'tx-7f3a:orders'; -- durable "yes" vote; locks are still held
-- Participant 2 (billing DB)
BEGIN;
INSERT INTO charges (id, order_id, cents) VALUES ('ch-91', 'o-55', 1299);
PREPARE TRANSACTION 'tx-7f3a:billing';
-- Coordinator writes its decision durably FIRST, then:
COMMIT PREPARED 'tx-7f3a:orders';
COMMIT PREPARED 'tx-7f3a:billing';
-- Recovery sweep after a coordinator crash: anything prepared and old is in doubt
SELECT gid, prepared, owner, database
FROM pg_prepared_xacts
WHERE prepared < now() - interval '5 minutes';
-- Resolve ONLY from the coordinator's decision log: COMMIT PREPARED or ROLLBACK PREPAREDThree production lessons follow from this. First, a prepared transaction survives restarts and keeps its row locks, and it also holds back the oldest transaction horizon, so vacuum cannot clean up dead rows behind it. A forgotten prepared transaction is a slow-motion outage. Second, resolution must come from the coordinator's durable decision log, never from guessing; if you cannot find the decision, escalate to a human rather than picking an outcome. Third, the global transaction id should encode the coordinator's identity and transaction id, as 'tx-7f3a:orders' does, so a recovery sweep can map every in-doubt branch back to a decision record.
Monitor the count and age of prepared transactions on every participant and alert on any older than a few minutes.
How distributed databases make commit cheap
Distributed SQL databases still run an atomic commit protocol, but they remove the two weaknesses of classic 2PC. Each participant is a replicated group that agrees through consensus, as described in Raft in depth, so a participant does not disappear when one machine does. And the commit decision is stored as data inside the database, where any node can read it, rather than in a separate coordinator's log.
Percolator, Google's incremental indexing system built on Bigtable, showed the pattern that TiDB and others follow. A transaction takes a start timestamp from a timestamp oracle and writes locks on every row it changes. One lock is the primary; every other lock points at it. Commit takes a commit timestamp and atomically replaces the primary lock with a commit record. That single-row write is the commit point. Secondary locks are cleaned up afterwards, and a reader that finds a stale secondary lock checks the primary to learn whether to roll it forward or back. Any client can finish a crashed transaction, so nothing waits for a coordinator.
Parallel commits in CockroachDB cut the latency further. Instead of writing intents, waiting, and then writing a committed transaction record, the coordinator writes the record in a STAGING state that lists the in-flight writes, at the same time as the writes themselves. The transaction is committed as soon as every listed write and the staging record have been replicated, so the client gets its answer after one round of consensus. The record is flipped to COMMITTED in the background. A reader that meets a STAGING record checks whether all listed writes succeeded and can decide the outcome itself.
Spanner runs two-phase commit across Paxos groups and adds external consistency with TrueTime, a clock API that returns an interval guaranteed to contain true time. The coordinator picks a commit timestamp no earlier than the latest possible current time, then waits until that timestamp is certainly in the past before making the commit visible. This commit wait is roughly twice the clock uncertainty. Systems without tightly bounded clocks use hybrid logical clocks instead, described in hybrid logical clocks, and accept uncertainty-driven read restarts.
Inside such a database, multi-shard transactions are correct and reasonably fast, but contention causes retries.
Retries are part of the contract
Serializable databases abort transactions that would violate isolation and return SQLSTATE 40001, serialization failure. That is not an error to log and forget; it tells the client to run the transaction again. Wrap every transaction in a retry loop with jittered backoff, and keep external side effects, such as HTTP calls and emails, out of the transaction body, because the body may run several times.
import random, time
from psycopg import errors
RETRYABLE = (errors.SerializationFailure, errors.DeadlockDetected) # SQLSTATE 40001, 40P01
def run_txn(pool, body, max_attempts=6):
for attempt in range(max_attempts):
with pool.connection() as conn:
try:
with conn.transaction():
return body(conn) # body must be free of external side effects
except RETRYABLE:
if attempt == max_attempts - 1:
raise
time.sleep(min(1.0, 0.01 * 2 ** attempt) * random.random()) # full jitterTrack the retry rate per transaction type; a rising rate usually means a hot row, best fixed in the data model.
Across services: a worked transfer
Now suppose accounts live in a Ledger service and user wallets in a separate Wallet service, each with its own database. No shared transaction is possible, and XA between teams is a coupling nobody wants. The goal is still that a transfer is never half-applied from the user's point of view, and that a retried request never debits twice.
The design combines three local mechanisms. The client sends an idempotency key with the request. The Ledger service claims the key, debits the account and writes an outbox event, all in one local transaction. A relay publishes the outbox event, and the Wallet service applies the credit idempotently, keyed by the same transfer id. If the credit is rejected, for example because the wallet is closed, Wallet emits CreditRejected and Ledger runs a compensating refund, which is a one-step saga.
# Transfer endpoint: idempotency key + local transaction + outbox, no cross-service lock.
import json, psycopg
class Conflict(Exception): pass # map to HTTP 409; client retries later
class InsufficientFunds(Exception): pass # map to HTTP 422
def transfer(conn, idem_key: str, src: str, dst: str, cents: int) -> dict:
with conn.transaction():
cur = conn.cursor()
# 1. Claim the key. Unique constraint makes the claim atomic.
cur.execute(
"INSERT INTO idempotency (key, status) VALUES (%s, 'started') "
"ON CONFLICT (key) DO NOTHING RETURNING key", (idem_key,))
if cur.fetchone() is None:
cur.execute("SELECT status, response FROM idempotency WHERE key = %s", (idem_key,))
status, response = cur.fetchone()
if status == "done":
return response # replay: same answer, no second debit
raise Conflict("request in progress") # concurrent duplicate
# 2. Local business change.
cur.execute(
"UPDATE accounts SET balance = balance - %s WHERE id = %s AND balance >= %s",
(cents, src, cents))
if cur.rowcount != 1:
raise InsufficientFunds(src)
# 3. Intent for the other service, in the SAME transaction.
event = {"type": "CreditRequested", "transfer": idem_key, "account": dst, "cents": cents}
cur.execute("INSERT INTO outbox (id, topic, payload) VALUES (%s, 'wallet', %s)",
(idem_key, json.dumps(event)))
response = {"transfer": idem_key, "state": "debited"}
cur.execute("UPDATE idempotency SET status = 'done', response = %s WHERE key = %s",
(json.dumps(response), idem_key))
return responseWalk through the failures. If the client times out and retries, the second request hits the idempotency row and gets the stored response, so there is no second debit. If Ledger crashes after the commit but before publishing, the relay publishes the outbox row when it restarts. If the relay publishes twice, Wallet's own processed-events table drops the duplicate. If Wallet is down for an hour, the event waits in the broker and the transfer shows as 'debited, crediting' in the meantime. The user-facing state machine must show that honestly, because this design gives eventual atomicity, not isolation.
Failure modes
| Symptom | Cause | Fix |
|---|---|---|
| Lock waits pile up on one database | In-doubt prepared transaction after coordinator crash. | Alert on prepared-transaction age; resolve from the decision log. |
| Double charge after a timeout | Retry without idempotency key, or key checked outside the transaction. | Claim the key with a unique constraint inside the same transaction. |
| Event published but database rolled back | Publishing to the broker inside the transaction body (dual write). | Outbox row in the same transaction; relay publishes after commit. |
| Retry storm under load | Hot row plus immediate retries. | Jittered backoff, retry budget, and split the hot row. |
| Email sent three times | Side effect inside a retried transaction body. | Move side effects to the outbox. |
| Stale lock or lease owner writes after losing ownership | Pause longer than the lease. | Fencing tokens checked by the storage layer; see fencing tokens. |
Operating it
- Measure the in-doubt window: prepared-transaction count and age, outbox lag (oldest unpublished row), saga instances stuck in a non-terminal state.
- Reconcile: run a daily job that compares both sides of every cross-service invariant, for example that every debit has a matching credit or refund, and pages when a mismatch is older than the maximum expected delay.
- Expire idempotency keys deliberately: keep them longer than the longest client retry horizon, and document that horizon in the API.
- Test crashes, not just happy paths: kill the process between each pair of steps in staging and verify that recovery converges.
Trade-offs
Co-location gives the strongest guarantee for the least effort but constrains your service boundaries. A distributed SQL database gives real transactions at the cost of commit latency, which grows with the distance between replicas, and a retry-aware client. XA gives atomicity across products but couples their availability and turns every coordinator bug into locked rows. Sagas and outboxes keep services independent but move isolation and reconciliation into application code, which you must test and monitor.
What to do next
- List every operation that changes data in more than one place, and write down the invariant each one must preserve.
- For each, try rung 1 first: can a partition key or ownership change make it local?
- If you run XA, add alerts on prepared-transaction age and write the recovery runbook now.
- Wrap every database transaction in a retry loop that handles SQLSTATE 40001 with jitter, and move side effects out of the body.
- Add idempotency keys to every mutating API that clients may retry, claimed inside the business transaction.
- Replace any publish-inside-transaction code with an outbox.
- Build a reconciliation job for each cross-service invariant and alert on old mismatches.