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 #systemdesignIn 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) is an academic solution involving a coordinator and multiple participants. It operates in two phases:
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.
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.
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.