You’re booking a trip. Flight, hotel, car — three services, three databases, three teams. The flight books. The hotel books. The car rental returns a 500.
Now what?
In a single database this is a non-question. You opened a transaction, you roll it back, the partial work vanishes and nobody ever knew. Across three services there is no transaction to roll back. The flight is booked. Somebody owns that seat now, and it’s your customer, and they’re expecting a car.
Everything below is about that gap.
2PC: the protocol that actually gives you atomicity
Two-phase commit is the honest answer to “I want a real distributed transaction,” and it does work. A coordinator runs it in two rounds.
Prepare. The coordinator asks every participant: can you commit this? Each one does the work, makes it durable, takes whatever locks it needs, and answers yes or no. A yes is a promise — it means “I can commit this whenever you say, and I will not change my mind.” That promise is the whole basis of the protocol.
Commit. If every answer was yes, the coordinator tells everyone to commit. If any answer was no, it tells everyone to abort.
That gives you genuine all-or-nothing. It’s not a trick or an approximation, and if you have XA-capable resource managers inside one trust boundary — a couple of databases and a message broker in the same datacentre — 2PC is a reasonable, well-understood choice. People who say “never use 2PC” are overcorrecting.
The problem is what happens between the phases.
Every participant that voted yes is now holding locks and waiting. It cannot commit, because it hasn’t been told to. It cannot abort, because it promised not to. It is stuck, by design, until the coordinator speaks.
So if the coordinator dies after collecting the votes but before announcing the decision, every participant stays blocked, holding locks, indefinitely. They can’t even safely ask each other — a participant that voted yes has no idea whether some other participant voted no. The decision existed only in the coordinator’s head, and the coordinator is gone.
This is the blocking problem, and it isn’t an implementation flaw you can engineer around. It’s inherent: 2PC can’t make progress through the failure of a single, specific node at a single, specific moment. In practice you mitigate it — the coordinator writes its decision to a durable log before announcing it so a restarted coordinator can finish, which is the same instinct as write-ahead logging. That converts “blocked forever” into “blocked until the coordinator comes back.” Better. Still blocked.
And the cost shows up long before anything fails. You’re holding locks across a network round trip, which caps throughput at something like one transaction per lock per round trip. Over a WAN, with services owned by different teams on different release cycles, one slow participant becomes everyone’s latency. It’s a protocol that assumes participants are nearby, fast, and administered together. Microservices across regions are none of those.
Sagas: give up atomicity, keep moving
The saga takes the opposite trade. Instead of one atomic operation across three services, you run three separate local transactions — each committing immediately in its own database — plus a plan for what to do when a later one fails.
T1 book flight C1 cancel flight
T2 book hotel C2 cancel hotel
T3 book car C3 cancel car
Run T1, T2, T3 in order. If T3 fails, run C2 and then C1. Each Cn semantically undoes its Tn.
Nothing is locked across services. Every step commits and releases. If the car service is down for an hour, the flight booking isn’t holding anything hostage. That’s the win, and it’s a large one.
The price is stated plainly: there is no isolation. Between T1 and T3, the world can see a state your business logic never intended — a flight booked for a trip that doesn’t exist yet. Another process reading that data sees it as real, because it is real; it’s committed. Anyone who reads it and acts on it has acted on a half-finished transaction, and if the saga later compensates, that reader’s decision was based on something that got undone.
The usual mitigations are a status field (PENDING until the saga completes, and readers are expected to respect it) or a semantic lock — marking the record as reserved so other sagas skip it. Both work. Both mean the isolation you lost is now your application’s problem, expressed in your domain’s vocabulary. That’s not free, and pretending otherwise is how sagas get a reputation for surprising people.
Compensation is not rollback, and the difference will bite you
This is the part that matters most, and it’s the part that gets glossed.
A database rollback erases. The change was never visible outside the transaction; afterwards there is no evidence it happened.
A compensation is a new action, taken after the original one was already visible. It doesn’t erase anything. It adds a second event that counteracts the first in business terms.
Refund a payment and the customer saw the charge. Their statement will show both the charge and the refund. Their bank may have charged an overdraft fee in between, and your refund doesn’t touch that. If it was a foreign-currency card, they may get back less than they paid because the rate moved. You have not restored the prior state. You’ve made a compensating entry, which is exactly what accountants have done for centuries — and notably, accountants never claimed it was the same as the transaction not happening.
So the design question for every step isn’t “can I undo this?” It’s “what is the business-meaningful counter-action, and is it good enough?” Sometimes the answer is that there isn’t one:
Sending an email. There’s no unsend. The best compensation is a second email saying disregard the first, which is worse than not sending the first.
Firing a missile, moving a robot arm, dispensing a reagent. Physical actions don’t have compensations in any useful sense — the reason I’ve argued that agents operating physical hardware break the control model most agent platforms assume.
Transferring money to an external bank. Now the compensation requires the recipient’s cooperation, which you do not control.
The standard move — and it’s a good one — is to order the steps so the irreversible ones go last. Do everything reversible first. Put the step with no compensation at the end, where if it fails you can still unwind everything before it, and if it succeeds you’re done. If you have two steps with no compensation, you have a design problem that sagas can’t solve, and it’s worth finding out at design time rather than at 2am.
A related rule: compensations must be idempotent, and so must the forward steps. The saga coordinator will retry. It will crash after a step succeeded but before it recorded that fact, restart, and run the step again. Without idempotency keys on both directions you get double bookings and double refunds, which is the failure mode that makes customers call their bank. Compensations also need to tolerate running against a step that never actually happened — a cancel for a booking that was never made should succeed quietly, not throw.
Who drives it
Two options, and the trade is readability versus coupling.
Orchestration puts a single component in charge. It calls T1, waits, calls T2, waits, and on failure walks the compensations backwards. The sequence lives in one file. You can read it. During an incident you can query the orchestrator and ask what state a given saga is in, and get an answer.
Choreography has no central component. Each service emits an event, and other services react. Booking the flight emits FlightBooked; the hotel service listens and books; and so on. Nothing is coupled to a coordinator.
The pitch for choreography is loose coupling, and it’s real. The cost is that the saga doesn’t exist anywhere. There’s no artifact describing the flow — it’s an emergent property of which services happen to subscribe to which events. Debugging means reconstructing the sequence from logs across every participant, and the people doing that reconstruction are usually tired and under pressure.
I default to orchestration for anything with more than about three steps or anything a human will have to debug. The central component people worry about is less risky than it sounds — it’s a state machine with a durable log, not a 2PC coordinator, and if it dies the saga just pauses and resumes when it’s back. Nobody is holding locks.
Choreography earns its place where the downstream reactions are genuinely independent and nobody needs to reason about the whole sequence. “Order placed, so send a receipt and update the recommendation model” is fine. “Order placed, so do seven things that must all succeed or all unwind” is not.
Picking
Honestly, the decision is usually made for you.
If every participant is a resource manager you administer, in one datacentre, and the transaction is short — 2PC is fine and simpler than people pretend. A database and a message broker committing together is the classic case.
If participants are services owned by other teams, across regions, behind APIs you don’t control, or on release schedules you can’t coordinate — 2PC isn’t available in any practical sense, and a saga is what you’re doing whether you’ve named it or not. The common failure here isn’t choosing the wrong pattern. It’s writing an implicit saga: a chain of service calls with ad-hoc error handling, no compensation for step two when step three fails, and no record of which sagas are half-finished. That’s a saga with the reliability parts left out. Naming the pattern is most of the fix, because once you’ve named it, “what’s the compensation for this step?” becomes a question somebody asks in review.
Before you build either, make sure you have the three things both depend on: idempotent operations, a durable record of where each transaction has got to, and timeouts with deadline propagation so a stuck step eventually becomes a failed step rather than an open-ended wait. Without those, 2PC blocks and sagas quietly strand themselves halfway, and the second one is harder to notice.
And ask the question that avoids all of this: does this actually need to span three services? A surprising number of distributed transactions exist because of a service boundary someone drew in a diagram three years ago. Moving the boundary so the operation is local to one service is allowed, and it is cheaper than every pattern on this page.
Related: Idempotency Keys: The Thing That Makes a Retry Safe · Write-Ahead Logging · Queues and Message Brokers
Comments