8.5 Distributed Transactions
Single-node atomicity works because of one physical fact:
- Make the transaction’s writes durable, typically in a WAL.
- Append a commit record to the log on disk.
“It is a single device — the controller of one particular disk drive, attached to one particular node — that makes the commit atomic.”
- Crash before the commit record is written → rolled back on recovery.
- Crash after → the transaction is committed.
Distributed: you cannot just send commit to every node.
“Once a transaction has been committed on one node, it cannot be retracted if it later turns out it was aborted on another node — because once committed, the data becomes visible to other transactions under read-committed or stronger isolation. If user 1’s transaction were later aborted, user 2’s transaction would have to be reverted as well, since it was based on data retroactively declared not to have existed.”
Ensuring the nodes either ALL commit or ALL abort is the ATOMIC COMMITMENT PROBLEM.
(Note: concurrency control in distributed transactions is broadly similar to single-node — serial execution on sharded databases, 2PL works distributed, and there are distributed serializability checkers for SSI. Achieving ATOMICITY is the new challenge.)
5.1 Two-Phase Commit (2PC)
The coordinator writes the decision to its transaction log on disk. That write — not any message — is the commit point.
If a commit request fails or times out, the coordinator must retry forever — there is no going back.
(The marriage analogy: the officiant asks each partner individually; after receiving both "I do"s, the couple is pronounced married — the transaction is committed — and the fact is broadcast. If either does not say yes, the ceremony is aborted.)
The system of promises — six steps
- The application requests a globally unique transaction ID from the coordinator
- It begins a single-node transaction on each participant, attaching the global transaction ID. All reads and writes happen in these single-node transactions. If anything goes wrong at this stage, the coordinator or any participant can abort
- When ready to commit, the coordinator sends
prepareto all participants, tagged with the global ID. If any request fails or times out, the coordinator sendsabortfor that ID to all participants - On receiving
prepare, a participant makes sure it CAN DEFINITELY COMMIT UNDER ALL CIRCUMSTANCES. This includes writing all transaction data to disk (A CRASH, A POWER FAILURE, OR RUNNING OUT OF DISK SPACE IS NOT AN ACCEPTABLE EXCUSE FOR REFUSING TO COMMIT LATER) and checking for conflicts/constraint violations. By replying yes, the node PROMISES to commit without error if requested — it SURRENDERS THE RIGHT TO ABORT, but without actually committing - The coordinator makes a definitive decision (commit only if all voted yes) and MUST WRITE THAT DECISION TO ITS TRANSACTION LOG ON DISK. ← the COMMIT POINT
- Once written, send commit or abort to all participants. If this fails or times out, THE COORDINATOR MUST RETRY FOREVER. There is no more going back. If a participant crashed meanwhile, the transaction will be committed when it recovers — SINCE IT VOTED YES, IT CANNOT REFUSE TO COMMIT WHEN IT RECOVERS
Two crucial POINTS OF NO RETURN: (1) when a participant votes yes, it promises it will definitely be able to commit later (though the coordinator may still choose to abort); (2) once the coordinator decides, that decision is IRREVOCABLE. Those promises ensure atomicity.
Single-node atomic commit lumps these two events into one: writing the commit record to the transaction log.
(Continuing the analogy: if you faint after saying "I do" and don't hear the pronouncement, that doesn't change the fact that the transaction was committed. When you recover, you can find out whether you are married by QUERYING THE OFFICIANT for the status of your global transaction ID — or wait for the officiant's next RETRY, since retries continued throughout your unconsciousness.)
Coordinator failure — the in-doubt problem
DB2 has committed. DB1 is now “in doubt”, or “uncertain”: it cannot unilaterally abort, because that would be inconsistent with DB2, which committed; and it cannot unilaterally commit, because another participant may have aborted. A timeout does not help.
“In principle, the participants could communicate among themselves to find out how each voted and come to an agreement — but that is not part of the 2PC protocol.” The only way 2PC can complete is by waiting for the coordinator to recover.
On recovery, the coordinator reads its transaction log; any transaction without a commit record is aborted.
Thus THE COMMIT POINT OF 2PC COMES DOWN TO A REGULAR SINGLE-NODE ATOMIC COMMIT ON THE COORDINATOR.
If the coordinator's DISK FAILS and its log is lost, the system HAS NO WAY TO AUTOMATICALLY RECOVER — only an administrator manually committing or aborting. And if only the most recent part of the log is lost, THE RECOVERING COORDINATOR MAY BELIEVE ALREADY-COMMITTED TRANSACTIONS WERE NOT COMMITTED AND TRY TO ABORT THEM, VIOLATING ATOMICITY.
Three-phase commit (3PC) was proposed to make atomic commit nonblocking. However, 3PC ASSUMES a network with BOUNDED DELAY and nodes with BOUNDED RESPONSE TIMES; in most practical systems with unbounded network delay and process pauses (Ch 9), 3PC CANNOT GUARANTEE ATOMICITY.
A better solution in practice is to REPLACE THE SINGLE-NODE COORDINATOR WITH A FAULT-TOLERANT CONSENSUS PROTOCOL (Ch 10).
5.2 Two very different kinds of distributed transaction
Distributed transactions have a MIXED REPUTATION: seen as providing an important safety guarantee hard to achieve otherwise, but criticized for causing operational problems, killing performance, and PROMISING MORE THAN THEY CAN DELIVER. Many cloud services choose not to implement them. (Much of the cost is the ADDITIONAL
fsyncOPERATIONS required for crash recovery and the ADDITIONAL NETWORK ROUND TRIPS.)
But two things get conflated:
| Database-internal | Heterogeneous | |
|---|---|---|
| Participants | All nodes run THE SAME database software | Two or more TECHNOLOGIES — databases from different vendors, or non-database systems like message brokers |
| Examples | YugabyteDB, TiDB, FoundationDB, Spanner, VoltDB, Cassandra, MySQL Cluster NDB, Kafka | XA transactions |
| Freedom | Doesn't have to be compatible with anything else → can use ANY protocol and apply technology-specific optimizations. CAN OFTEN WORK QUITE WELL | A LOT MORE CHALLENGING |
5.3 Exactly-once message processing via 2PC
A message from a queue can be acknowledged as processed IF AND ONLY IF the database transaction for processing it was successfully committed — by atomically committing the message acknowledgment and the database writes in a single transaction. With distributed transaction support this is possible even if the broker and the database are two unrelated technologies on different machines.
If either fails, both are aborted, so the broker may safely REDELIVER later. The abort discards any side effects of the partially completed transaction. This is EXACTLY-ONCE SEMANTICS.
⚠️ Only possible if ALL affected systems can use the SAME atomic commit protocol. If a side effect is sending an email and the email server does not support 2PC, THE EMAIL COULD BE SENT TWO OR MORE TIMES.
5.4 XA transactions
X/Open XA (eXtended Architecture), introduced 1991. Supported by PostgreSQL, MySQL, Db2, SQL Server, Oracle, and brokers ActiveMQ, HornetQ, MSMQ, IBM MQ.
XA is NOT A NETWORK PROTOCOL — it is merely a C API for interfacing with a transaction coordinator. In Java: JTA, supported by JDBC drivers and JMS broker drivers.
The architecture — and its fatal shape:
The coordinator reaches each system through its driver, which exposes callbacks for prepare, commit and abort. “The database server cannot contact the coordinator directly, since all communication must go via its client library.”
So if the application process crashes, or the machine dies, the coordinator goes with it. Any prepared-but-uncommitted participants are stuck in doubt. That server must be restarted and the coordinator library must read its log to recover.
Holding locks while in doubt — why in-doubt is so damaging
Database transactions acquire ROW-LEVEL EXCLUSIVE LOCKS on rows they modify. Under 2PL serializable, also SHARED LOCKS on rows they read. THE DATABASE CANNOT RELEASE THOSE LOCKS UNTIL THE TRANSACTION COMMITS OR ABORTS.
Therefore a transaction must HOLD ITS LOCKS THROUGHOUT THE TIME IT IS IN DOUBT. If the coordinator crashed and takes 20 minutes to restart, those locks are held for 20 MINUTES. IF THE COORDINATOR'S LOG IS ENTIRELY LOST, THOSE LOCKS ARE HELD FOREVER — or until manually resolved.
While held, no other transaction can modify those rows; depending on isolation level, others may even be blocked from READING them. THIS CAN CAUSE LARGE PARTS OF YOUR APPLICATION TO BECOME UNAVAILABLE.
Recovering from coordinator failure
In practice, ORPHANED IN-DOUBT TRANSACTIONS DO OCCUR — transactions whose outcome the coordinator cannot decide (log lost or corrupted by a software bug). These cannot be resolved automatically, so THEY SIT FOREVER IN THE DATABASE, HOLDING LOCKS AND BLOCKING OTHER TRANSACTIONS.
Even rebooting your database servers will not fix this, since a correct 2PC implementation MUST PRESERVE THE LOCKS OF AN IN-DOUBT TRANSACTION EVEN ACROSS RESTARTS (otherwise it risks violating atomicity). It's a sticky situation.
The only way out: an administrator examines the participants of each in-doubt transaction, determines whether any has already committed or aborted, and applies the same outcome to the others. This requires a lot of manual effort, most likely under high stress and time pressure during a serious production outage (otherwise, why would the coordinator be in such a bad state?).
Heuristic decisions — the emergency escape hatch letting a participant unilaterally decide without the coordinator:
To be clear, "heuristic" here is a EUPHEMISM FOR PROBABLY BREAKING ATOMICITY, since the heuristic decision violates the system of promises in 2PC. Intended only for catastrophic situations, not regular use.
The four fundamental problems with XA
- A single-node coordinator is a SINGLE POINT OF FAILURE for the entire system
- Making it part of the application server is problematic — the coordinator's logs on local disk become crucial durable system state, as important as the databases themselves
- Even a replicated coordinator wouldn't fix it: XA provides NO WAY for coordinator and participants to communicate DIRECTLY — only via the application code and drivers. SO THE APPLICATION CODE WOULD BE THE SINGLE POINT OF FAILURE. Solving this would require totally redesigning how application code is run to make it replicated or restartable — perhaps similar to DURABLE EXECUTION (Ch 5). However, no tools seem to take this approach in practice
- XA is a LOWEST COMMON DENOMINATOR, because it must be compatible with a wide range of systems. It cannot detect DEADLOCKS across different systems (that would require a standardized protocol for exchanging lock-wait information), and it DOES NOT WORK WITH SSI (that would require a protocol for identifying conflicts across systems)
5.5 Database-internal distributed transactions — why they're fine
NewSQL databases use 2PC for cross-shard atomicity yet don't suffer XA's problems, because they don't need to interface with other technologies — they avoid the lowest-common-denominator trap.
The four fixes:
- REPLICATE THE COORDINATOR, with automatic failover if the primary crashes
- Let the coordinator and data shards COMMUNICATE DIRECTLY, without intermediary application code
- REPLICATE THE PARTICIPATING SHARDS, reducing the risk of aborting because of a fault in one shard
- COUPLE the atomic commitment protocol WITH a distributed CONCURRENCY CONTROL protocol supporting deadlock detection and consistent reads across shards
Consensus algorithms replicate both coordinator and shards (Ch 10) — tolerating faults by automatically failing over WITHOUT HUMAN INTERVENTION while continuing to guarantee strong consistency. Both snapshot isolation and SSI are possible across shards.
5.6 Exactly-once WITHOUT distributed transactions
You don't actually need distributed transactions to achieve exactly-once semantics.
- Every message has a unique ID, and the database has a table of processed message IDs. On receiving a message, begin a database transaction and check the ID. Already present? Acknowledge to the broker and drop the message.
- Not present? Insert the ID, process the message with any additional writes going in the same transaction, and commit.
- Once committed, acknowledge the message to the broker.
- Once acknowledged, delete the message ID from the database, in a separate transaction.
| Crash | Outcome |
|---|---|
| Before step 2 commits | Transaction aborted; the broker retries; clean. |
| After step 2 but before step 3 | The broker retries; the retry sees the ID and drops it. |
| After step 3 but before step 4 | An old message ID lies around — harmless besides a little storage space. |
| A retry racing the original | A uniqueness constraint on the message-ID table prevents two concurrent transactions inserting the same ID. |
Achieving exactly-once processing requires ONLY TRANSACTIONS WITHIN THE DATABASE — atomicity across database and message broker is NOT NECESSARY. Recording the message ID makes the processing IDEMPOTENT, so it can be safely retried without duplicating side effects. A similar approach is used in Kafka Streams (Ch 12).
That said, internal distributed transactions are still useful for the SCALABILITY of such patterns — e.g. message IDs on one shard and the main data on other shards, with atomicity across those shards.