Distributed Transactions

Distributed Transactions

Definition: Mechanisms designed to enforce transactional consistency across operations that span multiple independent services or database nodes, where no single database engine can natively guarantee atomicity.

How It Works

  • Two-Phase Commit (2PC): a coordinator sends a Prepare message to every participant. Each participant locks its resources, does the work, and votes YES (ready to commit) or NO. If all vote YES, the coordinator sends Commit; if any votes NO, or times out, it sends Abort. All participants then apply the outcome and release locks.
  • The actual message sequence, phase by phase:
    1. Phase 1, Voting (Prepare) Phase: coordinator → all participants: PREPARE. Each participant performs the work locally (writes to its log, acquires locks) without committing, then replies VOTE-COMMIT or VOTE-ABORT. A participant that votes VOTE-COMMIT has made a durable promise: it must be able to commit later even if it crashes and restarts in between.
    2. Coordinator collects all votes. If every participant voted VOTE-COMMIT, it durably logs GLOBAL-COMMIT before sending anything else, this log write is what makes the coordinator’s decision recoverable after a crash. If any participant voted VOTE-ABORT (or didn’t respond in time), it logs GLOBAL-ABORT instead.
    3. Phase 2, Commit Phase: coordinator → all participants: GLOBAL-COMMIT or GLOBAL-ABORT. Each participant applies the decision (commits or rolls back), releases its locks, and sends an ACK back to the coordinator.
    4. Coordinator waits for all ACKs before considering the transaction fully closed and discarding its own log record for it.
  • The blocking failure mode lives specifically between step 1 and step 3: a participant that has voted VOTE-COMMIT but hasn’t yet received GLOBAL-COMMIT/GLOBAL-ABORT is “in doubt”, it cannot unilaterally decide either way, since the coordinator might have already told other participants to commit.
  • 2PC is blocking by construction: after voting YES, a participant must hold its locks until it hears back from the coordinator. If the coordinator crashes after collecting votes but before broadcasting the outcome, participants are stuck holding locks indefinitely, unable to unilaterally decide commit or abort.
  • Three-Phase Commit (3PC) adds a “pre-commit” phase and timeouts to avoid indefinite blocking on coordinator failure, but it assumes bounded network delay (no partitions), which real networks don’t guarantee, so it’s rarely used in practice.
  • Saga pattern: breaks a transaction into a sequence of local transactions, each committed independently in its own service/database. If a later step fails, previously completed steps are undone by explicit compensating transactions run in reverse order, rather than a coordinator-enforced rollback.
  • Choreography-based sagas: each service publishes an event when its step completes; the next service reacts to that event and performs its own step. No central coordinator, fully decentralized, but the overall transaction flow is implicit and harder to trace across services.
  • Orchestration-based sagas: a central orchestrator explicitly calls each step and, on failure, explicitly calls the corresponding compensating actions. Easier to reason about and debug, at the cost of a central component that now knows about every service’s workflow.
  • Outbox pattern: a common building block for sagas, where a service writes its local state change and an event describing it in the same local transaction (to the same database), then a separate process reliably publishes that event. This avoids the classic “database committed but the message never got published” failure mode of dual writes.
  • Compensating transactions are not true rollbacks: they are new forward-moving actions that semantically undo a prior effect (e.g., “refund the charge” rather than “pretend the charge never happened”), because in most cases the original effect (an email sent, a payment captured) can’t be literally erased.
  • Semantic locks: because a saga’s intermediate state is visible to other transactions, some implementations mark in-progress records with a status flag (e.g., “reserved, pending payment”) so concurrent transactions can react appropriately instead of treating a not-yet-finalized state as final.
  • Try-Confirm/Cancel (TCC): a variant closer to 2PC in structure but implemented at the application level: each participant exposes a “try” (tentatively reserve resources), “confirm” (finalize), and “cancel” (release) operation, giving explicit two-phase semantics without a blocking, lock-holding coordinator.

Worked Example: Timing a 2PC Transaction vs a Saga

A checkout touching 3 services, each with 20ms local processing time and 10ms one-way network latency to the coordinator/orchestrator:

  • 2PC: Prepare round trip (10ms there, 20ms process, 10ms back = 40ms, done in parallel across participants so bounded by the slowest) + coordinator decision + Commit round trip (another ~40ms) = roughly 80ms minimum, plus every participant holds its locks for the full duration, meaning throughput on those specific rows is capped by how many of these 80ms round trips can be pipelined.
  • Saga (orchestrated): each step commits locally and returns before the orchestrator calls the next one: 3 sequential steps × (10ms there + 20ms process + 10ms back) = roughly 120ms end to end, slightly slower wall-clock for this happy path, but no service ever holds a lock waiting on another service’s network round trip, so throughput per service is bounded only by that service’s own capacity, not by the slowest participant in the whole chain.
  • The saga’s real advantage isn’t shown by this single-request timing, it’s that under load, 2PC’s held locks create contention that compounds, while the saga’s independent local commits don’t, which is why the gap widens dramatically at scale even though a single saga transaction can look no faster (or even slower) than 2PC in isolation.

Trade-offs

  • 2PC gives real atomicity (all participants commit or all abort, and intermediate states are never visible) at the cost of blocking behavior, a coordinator single point of failure, and locks held across a network round trip, all of which hurt throughput and availability.
  • Sagas avoid 2PC’s blocking and availability problems, since each local transaction commits independently and quickly, but they give up true atomicity: intermediate states are visible to other transactions while a saga is in progress, and a failure partway through leaves the system in a state that must be explicitly compensated, not automatically rolled back.
  • Choreography sagas avoid a central point of failure or bottleneck but trade away a single place to read “what does this transaction actually do,” making debugging and onboarding harder as the number of steps grows.
  • Orchestration sagas centralize workflow logic for easier reasoning and observability, but reintroduce a coordinator dependency, though a much lighter one than 2PC’s, since the orchestrator doesn’t hold locks across the network the way a 2PC coordinator does.
  • Every compensating action must itself be designed, and not every action is cleanly compensatable (an email can be “corrected” with a follow-up, but not unsent), which pushes real design effort into the failure path, not just the happy path.

Why It Matters

  • It’s what makes correctness possible for business processes that inherently span multiple services or databases (checkout, order fulfillment, financial transfers), which is the norm in a microservice architecture, not the exception.
  • It maintains business-state consistency across decoupled services without requiring a single, giant, tightly-coupled database that would defeat the purpose of splitting into services in the first place.
  • The choice between 2PC and sagas is a direct proxy for a service’s tolerance for availability loss versus its tolerance for temporarily-visible intermediate states, a distinction every checkout, payments, or booking flow has to make explicitly.

Common Pitfalls

  • Using 2PC in a high-throughput, cloud-scale system without recognizing the latency and availability cost: every transaction now waits on the slowest participant and holds locks across a network round trip, which doesn’t scale the way independent local transactions do.
  • Designing a saga without a compensating action for every step that has a side effect, discovering only in production that some step (e.g., “email sent”) can’t actually be undone.
  • Forgetting that saga steps need to be idempotent, since retries (from timeouts, redelivery, or crash recovery) can cause a step or its compensation to run more than once.
  • Assuming a saga gives isolation. Other transactions can observe a saga’s intermediate, not-yet-fully-committed state (e.g., inventory reserved but payment not yet captured), which needs explicit handling (semantic locks, status flags) if that visibility would cause a business problem.
  • Building a 2PC coordinator without persisting its decision before crashing. If the coordinator’s outcome isn’t durably logged before a crash, recovering participants can be left permanently stuck, unable to determine whether to commit or abort.
  • Reaching for distributed transactions when the actual fix is a better data model, co-locating data that’s always modified together on the same shard or service removes the need for cross-service coordination entirely.
  • Skipping the outbox pattern and doing a “dual write” instead, writing to the local database and publishing an event as two separate operations, which leaves a window where one succeeds and the other doesn’t, with no transaction tying them together.
  • Assuming an orchestrator’s own failure is someone else’s problem. The orchestrator’s state (which step a saga is on) needs to be persisted and recoverable, or a crash mid-saga leaves steps that already ran with no record of needing compensation.

Comparison

Two-Phase Commit (2PC)Saga PatternSingle-Node ACID Transaction
AtomicityTrue, all-or-nothingNo, compensations undo forward-completed stepsTrue, native to the database engine
BlockingYes, participants hold locks awaiting coordinatorNo, each local transaction commits independentlyBriefly, within the single transaction
Availability impactHigh, coordinator failure can stall participantsLow, no cross-service locks heldLow, entirely local
Failure recoveryCoordinator must resolve in-doubt transactionsExplicit compensating transactionsStandard rollback
Typical useRare in modern cloud systems, some legacy/XA integrationsMicroservice workflows (checkout, order fulfillment)Anything within a single database
Locks held across networkYes, until coordinator resolvesNoNo, locks are local and brief
Failure mode if a participant is unreachableWhole transaction blocksOnly that step’s business flow is affected, others already committedN/A, no network involved

Real-World Scenario

An e-commerce checkout spans three services: inventory, payments, and shipping. Using a saga: 1) inventory service reserves stock and commits locally, 2) payments service charges the card and commits locally, 3) shipping service creates a shipping order and commits locally. If step 2 fails (card declined), the orchestrator calls a compensating action on inventory (release the reservation) and the saga ends in a defined “failed” state rather than a half-committed one. No cross-service lock was ever held; each step was fast and independently committed, at the cost of a brief window where inventory was reserved but no payment had yet succeeded.

Debugging Walkthrough: A Coordinator Crash Mid-Commit

  1. A legacy order-processing system uses XA-based 2PC across an orders database and a billing database. The coordinator sends PREPARE to both; both lock their rows, do the work, and reply VOTE-COMMIT.
  2. The coordinator process crashes (out-of-memory kill) immediately after receiving both votes but before it durably logs GLOBAL-COMMIT and before it sends the commit message to either participant.
  3. Both participants are now “in doubt”: they’ve voted VOTE-COMMIT, meaning they’re holding row locks and cannot unilaterally abort (the coordinator might have already told the other participant to commit, and aborting unilaterally would break atomicity), but they also haven’t been told to commit.
  4. On-call sees the symptom as a sudden spike in lock-wait time and blocked queries on both the orders and billing tables, transactions from other requests start timing out waiting for the same rows.
  5. The coordinator restarts. Because it never logged GLOBAL-COMMIT before crashing, its recovery routine treats the transaction as undecided and must ask each participant for its vote status, or, if it also lost that record, it may have to default to GLOBAL-ABORT as the safe choice, since it can’t prove a commit was ever finalized.
  6. If the coordinator had durably logged GLOBAL-COMMIT before crashing (the correct implementation), its recovery routine would replay that decision to both participants on restart, they’d commit and release their locks, and the blocked transactions would drain immediately once the coordinator comes back.
  7. Root cause of the incident: the coordinator’s decision log write and the crash race were close enough together that this class of bug went undetected until real production load hit it. The fix is ensuring the GLOBAL-COMMIT/GLOBAL-ABORT log write is fsync’d before any participant is notified, and that recovery always consults that log first, never re-deciding from scratch. This exact failure mode, and the operational cost of participants blocking indefinitely on a crashed coordinator, is the concrete reason most teams migrate off 2PC toward sagas for anything on the critical request path.

FAQ

Why isn’t 2PC popular in modern microservice architectures? Because it requires all participants to be available and responsive for the duration of the transaction, which conflicts directly with the availability and independent-deployability goals that led teams to microservices in the first place.

Do sagas provide isolation like ACID transactions do? No. Sagas provide eventual consistency across steps, not isolation; other actors can see intermediate states while a saga is still executing.

What’s the difference between a saga’s compensating transaction and a database rollback? A rollback undoes an uncommitted change atomically and invisibly. A compensating transaction is a new, separately committed action that semantically reverses an already-committed effect, and it’s visible as its own operation.

Can a saga be partially compensated, leaving the system in a genuinely inconsistent state? Yes, if a compensating action itself fails (e.g., the refund API is down). This is why compensations typically need their own retry-with-backoff and, as a last resort, an alert for manual intervention, since a saga’s safety net has a failure mode too.

Is XA a form of 2PC? Yes. XA is a widely implemented standard protocol for 2PC across heterogeneous resource managers (multiple databases, message queues), and it inherits 2PC’s blocking behavior and coordinator dependency.

History

  • 2PC originated in distributed database research in the late 1970s and was standardized industry-wide as the XA specification in 1991, letting heterogeneous resource managers (databases, message queues) participate in a shared transaction coordinated by a transaction manager.
  • The Saga pattern was introduced in a 1987 paper by Hector Garcia-Molina and Kenneth Salem as a way to handle “long-lived transactions” without holding locks for their full duration, predating microservices by decades but becoming the standard answer to cross-service consistency once service-oriented and microservice architectures took hold in the 2010s.
  • The outbox pattern and event-driven choreography became common saga-implementation techniques alongside the rise of message brokers (Kafka, RabbitMQ) as the default microservice integration layer.

Common Interview Questions

  • Why is 2PC considered a poor fit for microservices? It requires all participants to hold locks and stay available for the duration of the transaction, directly conflicting with independent deployability and availability goals of microservice architecture.
  • What’s the core difference between choreography and orchestration sagas? Choreography has services reacting to each other’s events with no central coordinator; orchestration has a central process explicitly directing each step and its compensations.
  • How do you make saga steps safe to retry? By designing each step (and its compensation) to be idempotent, so redelivering the same message or retrying after a timeout doesn’t double-apply the effect.
  • What’s a concrete example of a step that’s hard to compensate? Sending an email or an SMS; the compensation is a follow-up message, not an undo, since the original action can’t be unsent.
  • When would you actually still reach for 2PC today? Narrow cases with a small, fixed set of tightly-coupled resource managers that must be atomic and where availability during the transaction window is not a primary concern, such as certain legacy enterprise integrations.

Example

E-commerce checkout saga: 1) reserve inventory, 2) charge the credit card, 3) create a shipping order. If step 2 fails, run a compensating transaction to unreserve the inventory rather than attempting a cross-service rollback. Banking systems moving money between accounts at different institutions typically use saga-like patterns (or asynchronous settlement) rather than 2PC, since blocking two entire banks’ systems on one transaction isn’t operationally viable.

Design Checklist

  • Can the transaction be avoided entirely by co-locating the data it touches on the same shard or service, removing the need for cross-service coordination?
  • If a saga is chosen, does every step with a side effect have a defined, tested compensating action?
  • Are saga steps and their compensations idempotent, safe to run twice if a retry or redelivery happens?
  • Who, or what, is responsible for detecting a stuck or partially-completed transaction and driving it to a terminal state?

Dig deeper