Menu
Dev.to #systemdesign·September 28, 2026

Distributed Transactions: Why Two-Phase Commit Fails and Sagas Succeed

This article explores the challenges of distributed transactions in microservice architectures, contrasting the theoretical appeal of Two-Phase Commit (2PC) with its practical drawbacks at scale. It argues that 2PC's blocking nature makes it unsuitable for independent services and instead advocates for the Saga pattern, which trades strong consistency for eventual consistency and higher availability through compensating actions. The discussion highlights crucial trade-offs in designing reliable distributed systems.

Read original on Dev.to #systemdesign

The Challenge of Distributed Transactions

In a monolithic application with a single database, ACID transactions provide atomicity and isolation, ensuring that a series of operations either all succeed or all fail, leaving the system in a consistent state. However, when systems scale by sharding databases or adopting microservices, each with its own database, this inherent transactional guarantee evaporates. A single logical operation, like processing an order, now spans multiple independent services and databases, leading to the problem of distributed transactions. Without a coordinated mechanism, partial failures become common, resulting in inconsistent states (e.g., a card charged but inventory not reserved).

Two-Phase Commit (2PC): The Textbook but Impractical Solution

Two-Phase Commit (2PC) is an academic solution involving a coordinator and multiple participants. It operates in two phases:

  1. Phase one (Prepare): The coordinator asks all participants if they can commit. Each participant performs the work, durably records the change, locks affected rows, and votes 'yes' or 'no'.
  2. Phase two (Commit): If all participants vote 'yes', the coordinator tells everyone to commit and release locks. If any vote 'no', or if the coordinator fails, it tells everyone to abort and release locks.
⚠️

Why 2PC Fails in Distributed Systems

2PC is a blocking protocol. If the coordinator fails *between* collecting votes and broadcasting the decision, participants remain in a 'prepared' state, holding locks indefinitely. This blocks other transactions and makes the system's availability dependent on the uptime of *all* participants and the coordinator. A single slow participant can also stall the entire transaction, degrading system performance. This tight coupling and sensitivity to failures make 2PC impractical across independent services at internet scale.

Sagas: The Industry Standard for Eventual Consistency

Real-world systems like Uber and Netflix employ the Saga pattern to manage distributed transactions. Sagas embrace eventual consistency, acknowledging that a system may be temporarily inconsistent but will always converge to a consistent state. Instead of one global transaction, a saga is a sequence of local transactions, each committed immediately by its respective service.

When a step in the saga fails, there's no automatic rollback like with 2PC. Instead, compensating actions are triggered. These are business-level undo operations (e.g., a refund for a charge, releasing reserved stock). This approach avoids blocking and improves availability, as no service holds locks waiting for another independent service.

Where 2PC is Still Used (and Why)

While generally avoided for inter-service communication, 2PC is alive and well *within* highly coupled, distributed systems like Google Spanner and YugabyteDB. In these cases, the coordinator and participants are components of a single, tightly managed system, often backed by robust consensus protocols (like Paxos or Raft) for fault tolerance. The complexity of 2PC is abstracted away from the application developer, illustrating that its viability depends heavily on tight integration and controlled failure domains.

distributed transactionstwo-phase commitsagaseventual consistencyACIDBASEmicroservicestransaction management

Comments

Loading comments...