Why this matters
Inside one database, transactions are a solved problem: BEGIN, do work, COMMIT, and atomicity is guaranteed. Across multiple databases or services, that guarantee evaporates — and "book the flight, charge the card, reserve the hotel" either all happens or the customer has a very bad day. Distributed transactions are how you get atomicity back, and every approach is a trade-off between correctness, availability, and your own sanity. This is the lesson where the CAP theorem and the microservices lesson collide at full speed.
Analogy: a wedding planner coordinating five vendors — florist, caterer, band, photographer, venue — for one Saturday. The planner calls each one: "Can you commit to June 14th?" Everyone says yes (phase 1: prepare). Then the planner calls back: "It's on — lock it in" (phase 2: commit). But if the photographer stops answering between the two calls, the planner is stuck: the florist is holding the date, the money is half-moved, and nobody knows whether the wedding is happening. That stuck state has a name — blocking — and it's the central drama of this lesson.
Video: Distributed Transactions in Microservices | 2PC, Failures & Compensation | Episode 15 — Avinash Explains Tech
Covers local vs distributed transactions, ACID limits across services, and the partial-failure problem that starts it all.
Two-phase commit (2PC): the classic
One coordinator drives the transaction across participants (databases/services):
sequenceDiagram
participant C as Coordinator
participant P1 as Participant 1 (orders DB)
participant P2 as Participant 2 (payments DB)
C->>P1: PREPARE (can you commit?)
C->>P2: PREPARE (can you commit?)
P1-->>C: YES (durable, locked)
P2-->>C: YES (durable, locked)
Note over C: All YES → decide COMMIT
C->>P1: COMMIT
C->>P2: COMMIT
P1-->>C: ACK
P2-->>C: ACK
Phase 1 (prepare): each participant does the work, writes it durably, takes the locks, and votes YES — or votes NO and the whole thing aborts. Phase 2 (commit/abort): the coordinator broadcasts the decision; participants apply it.
The numbers: a commit costs 2 round trips minimum (prepare + commit), and participants hold locks for the entire duration — including while waiting on a slow coordinator. Lock-hold time directly caps throughput: if a 2PC transaction holds locks for 50 ms, that's at most ~20 such transactions per second per contended row, no matter how much hardware you buy. 2PC buys atomicity with latency and contention.
Video: Two-Phase Commit (2PC) Explained: Pros, Cons & When to Use It — Escoding
Walks through the prepare and commit phases, why coordinators matter, and where 2PC makes sense.
The blocking problem
Here's the wedding planner's nightmare, formalized. If the coordinator crashes after participants voted YES but before sending the decision, participants are blocked: they hold locks, they can't unilaterally commit (maybe the coordinator decided abort), and they can't unilaterally abort (maybe it decided commit — another participant may have already committed). They wait. And wait. Locks held, throughput zero, until the coordinator recovers.
Three-phase commit adds a pre-commit round so participants can time out and abort safely — but it only handles coordinator crashes, not network partitions, and the extra round trip made it unpopular. The industry verdict: 2PC is fine within one database cluster (where the coordinator is really the primary and failures are fenced), and increasingly avoided across independent services — which is where sagas come in.
Interactive diagram: StepThrough (loads in the app)
Video: Three Phase Commit Explained: How Distributed Systems Stay Consistent | Reserve, Pre-Commit, Commit — Kartikeya Sharma
Shows how a coordinator crash leaves participants holding locks in an uncertain state — then how 3PC answers it.
Sagas: embracing the lack of atomicity
A saga breaks the distributed transaction into a sequence of local transactions, each with a compensating action that semantically undoes it:
flowchart LR
T1["Book flight<br/>↩ cancel flight"] --> T2["Charge card<br/>↩ refund card"]
T2 --> T3["Reserve hotel<br/>↩ cancel hotel"]
T3 --> OK["✓ All done"]
T2 -->|charge fails| C1["Compensate:<br/>cancel flight"]
If step 3 fails, the saga runs compensations for steps 1–2 in reverse. Notice what's missing: isolation. Between "flight booked" and "flight cancelled," other transactions see the intermediate state — the seat looks taken, then freed. Sagas give you eventual atomicity (all-or-nothing eventually), not ACID atomicity. For most business processes (orders, bookings, onboarding flows), that's exactly the right trade — and it's how orchestrated (a central saga coordinator, e.g., Temporal/Cadence) or choreographed (services react to each other's events, no center) sagas run the world's e-commerce.
Compensating actions are the hard part: "refund the card" is easy; "un-send the email" is impossible — you send a correction email instead. Design compensations to be idempotent (the saga might retry them) and accept that some effects are only semantically undone, never truly erased.
Video: Not Just Events: Developing Asynchronous Microservices • Chris Richardson • GOTO 2019 — GOTO Conferences
Richardson explains sagas as sequences of local transactions, covering orchestration and messaging for reliability.
The outbox pattern: transactions + messaging, safely
The most common distributed-transaction bug in microservices: the service updates its database and then publishes an event — and crashes between the two. Database says "order placed," no event ever fires, downstream services never hear about it. Or the reverse: event published, DB write rolled back — a ghost order haunts the system.
The transactional outbox fixes it with one local transaction:
sequenceDiagram
participant S as Service
participant DB as Database
participant R as Relay
participant Q as Message queue
S->>DB: BEGIN
S->>DB: UPDATE orders SET status='placed'
S->>DB: INSERT INTO outbox (event)
S->>DB: COMMIT (atomic — both or neither)
R->>DB: Poll outbox (unpublished events)
R->>Q: Publish event
R->>DB: Mark published
The business write and the event land in the same local transaction — atomicity restored. A separate relay process publishes outbox rows to the queue (at-least-once, so consumers must be idempotent — there's our old friend from Message Queues). Debezium (CDC-based) is the industrial-strength version: it tails the database's own replication log instead of polling a table.
Video: Stop writing to two systems. Write to one. — Metaphorically Speaking
A card-game metaphor shows why writing to the DB and the broker separately loses events — and how one outbox table fixes it.
Choosing your poison
| Approach | Atomicity | Availability during failures | Complexity |
|---|---|---|---|
| 2PC/XA | Strong (ACID) | Blocks on coordinator failure | Moderate (built into DBs) |
| Saga (orchestrated) | Eventual | High (each step independent) | High (compensations) |
| Saga (choreographed) | Eventual | Highest | Highest (debugging events) |
| Outbox + events | Per-service local | High | Low–moderate |
The pragmatic default for microservices: avoid distributed transactions where the domain allows it (redesign boundaries so one service owns the write), use the outbox for "update + notify," and reach for sagas when a multi-step business process truly spans services. Reserve 2PC for single-cluster cases where the database does it for you.
Video: Distributed Transactions Explained: 2 Phase Commit vs Saga Pattern — Hello Interview
Compares 2PC and sagas head-to-head, covers the outbox as the glue, and gives a practical verdict.
Takeaways
- 2PC gives true atomicity at the cost of 2 round trips, long lock holds, and blocking if the coordinator dies mid-protocol.
- Sagas trade isolation for availability: local transactions plus idempotent compensating actions, run forward on success and backward on failure.
- The outbox pattern makes "DB write + publish event" atomic via one local transaction plus a relay — the fix for the dual-write bug.
- Default to avoiding cross-service transactions; use outbox for notify-flows and sagas for multi-step processes.
Check your understanding
In two-phase commit, what does it mean for participants to be 'blocked'?
- Two transactions are deadlocking on the same row
- The coordinator crashed after YES votes but before the decision, so participants hold locks indefinitely, unable to commit or abort
- They are waiting for more participants to join
- The network is too slow for the prepare phase
What does a saga use instead of rollback?
- Compensating actions that semantically undo each completed step, run in reverse order
- Automatic retries until the step succeeds
- A global lock held for the saga's duration
- Database savepoints
What bug does the transactional outbox pattern fix?
- Slow queries blocking the event loop
- Cache stampedes on hot keys
- The dual-write problem: crashing between a database update and publishing its event
- Deadlocks between concurrent transactions
What consistency does a saga provide compared to ACID?
- No consistency at all — steps are independent
- Full ACID atomicity and isolation
- Stronger than ACID, with causal ordering
- Eventual atomicity — all-or-nothing eventually, but intermediate states are visible
Go deeper
Want to keep pulling this thread? These talks and tutorials go further than we did here:
- microXchg 2018 - Managing data consistency in a microservice architecture using Sagas — Chris Richardson (microservices.io), microXchg 2018. Why 2PC fails across services; sagas with compensations and transactional messaging.
- Gunnar Morling & Hans-Peter Grahsl about "Change Data Streaming Patterns in Distributed Systems" — Konfy. The outbox pattern and CDC-based saga orchestration, implemented live.
- Sean T. Allen on Life Beyond Distributed Transactions: An Apostate's Opinion — Sean T. Allen, Papers We Love San Francisco. Pat Helland's thesis on why distributed transactions can't scale and what replaces them.
- Distribute your microservices data with events, CQRS, and event sourcing — Red Hat Developer, DevNation Tech Talk (~36m). How to split data across services without distributed transactions — events, CQRS, and event sourcing with Kafka.
- Distributed Consensus and Data Replication strategies on the server — Gaurav Sen. 2PC and Sagas/quorum covered as replication and consistency strategies.
Sources & further reading
- Jim Gray, "Notes on Data Base Operating Systems" (1978) — the paper that formally characterized two-phase commit (the protocol itself was published earlier by Lampson & Sturgis).
- Hector Garcia-Molina & Kenneth Salem, "Sagas" (ACM SIGMOD 1987) — the original saga paper.
- Chris Richardson, microservices.io, "Transactional Outbox" — the pattern write-up with implementation options.
- Martin Kleppmann, Designing Data-Intensive Applications (O'Reilly, 2017), Ch. 9 — distributed transactions, 2PC limitations, and exactly-once semantics.
- Better Engineers, "A Crash Course on Microservices Architecture" (Substack) — a concise refresher on the Saga pattern alongside the gateway, BFF, CQRS, and event-sourcing patterns.