Distributed Transactions & Two-Phase Commit

11.A Promise That Spans Machines

M

In this chapter

We'll trace two-phase commit through GreenMart's own real order (deduct stock, create order, across two shards) as the both-or-neither protocol for an operation spanning multiple machines, its honest coordinator-dependency weakness, and clock problems and replication lag as the same honest theme showing up even without a dramatic failure.

10–12 min

The Problem in Real Life

A single order now has to touch two genuinely different machines at once: deduct stock from the inventory shard that owns this product, and create the order record on the shard that owns this customer. Act 4 already taught a conditional write keeping one item honest on one machine. This is a different, harder question: what if the deduction succeeds and the order creation fails, or the other way around?

Sarah doesn't want "probably fine." She wants a real, both-or-neither guarantee, across two machines that don't share anything except a network connection.

S

Both machines have to agree to do this together, or neither of them should do it at all.

Sarah

One Machine's Own Promise vs. A Promise Two Machines Both Keep

A distributed transaction spans machines

One logical operation that has to succeed or fail as a whole, across more than one machine with no shared memory.

Two-phase commit: prepare, then commit

Every participant agrees it can proceed first; only if all agree does the coordinator tell everyone to actually commit.

The coordinator is a real weak point

A coordinator failure mid-protocol can leave participants blocked, holding locks, genuinely unsure what to do.

Clocks and replication lag are honest, quiet problems

Machines' clocks aren't perfectly synchronized, and a write takes real, measurable time to reach every follower.

Distributed Transactions & Two-Phase Commit

A distributed transaction is exactly this: a single logical operation that has to succeed or fail as one unit, across more than one machine — genuinely harder than a single-machine transaction, because there's no one place that can just lock everything and decide alone.

Two-Phase Commit — GreenMart's Own Order, Traced Through

Phase 1: Prepare

Coordinator asks: inventory shard, can you deduct 1 unit? order shard, can you create this order?

each participant checks and locks, without committing yet

Both reply

Inventory shard: "yes, I can, and I've locked it." Order shard: "yes, I can, and I've locked it."

only if every participant said yes

Phase 2: Commit

Coordinator tells both: commit for real, now

if even one had said no, instead

Phase 2: Abort

Coordinator would tell both: abort — nothing happens anywhere

Two-phase commit (2PC) is the classic, real protocol for exactly this, walked through concretely above. Phase one (prepare): a coordinator asks every participant — the inventory shard, the order shard — "can you do your part of this, and are you willing to guarantee it, if asked?" Each one checks, locks what it needs, and answers yes or no, without actually committing yet. Phase two (commit or abort): if every participant said yes, the coordinator tells all of them to actually commit, for real. If even one said no, it tells all of them to abort — nothing happens anywhere, cleanly. GreenMart never ends up in the broken state where stock was deducted but no order exists, or an order exists for stock that was never actually reserved.

2PC genuinely works, and it has a real, honest weakness worth naming directly: it depends on the coordinator. If the coordinator itself fails after phase one but before telling everyone the final answer, participants are left holding a lock, genuinely unsure whether to commit or abort — blocked, potentially for a long time, waiting on a machine that might not come back. This is exactly why real consensus protocols (Raft, Paxos, from two chapters back) are often preferred for genuinely critical coordination — they're built from the ground up to keep making progress even when some machines are unreachable, rather than depending on one coordinator's own survival.

Clock/time problems and replication lag are two more, real, honest reasons distributed correctness is hard even without an outright failure. Different machines' clocks are never perfectly synchronized — "which write happened first" can't always be answered by just comparing timestamps across machines, the same real problem behind why this course's own simulator-backed Acts always used a relative or simulated clock rather than real wall-clock time. Replication lag is the honest gap between a write landing on a leader and that same write actually reaching every follower — a real, measurable delay, not an instant, and exactly the kind of gap that can turn a read against a lagging follower into a stale one, even with no partition happening at all.

Key Takeaway

A distributed transaction needs a real, both-or-neither guarantee across machines that don't share memory — two-phase commit is the classic mechanism for that, real and correct, with an honest weakness (a coordinator failure can leave participants blocked) that's exactly why critical coordination often leans on majority-quorum consensus instead. Clock problems and replication lag are the same honest theme in a quieter, everyday form: even without a dramatic failure, distributed correctness has to account for time and propagation delay being real, not instant.

Why This Matters

This is the real mechanism behind any operation that has to stay correct across more than one of GreenMart's own shards or regions — and its honest weakness (coordinator dependency) directly explains why the consensus-based approach from two chapters back is often the more resilient real-world choice for anything truly critical.

GreenMart now has a real, named protocol for a promise spanning multiple machines, traced through its own actual order-plus-inventory example, and an honest account of where it can still get stuck. What happens when a whole region — not just one link — genuinely goes down is exactly where the final chapter goes.

Next