Hello! In our previous lessons, we've explored how to scale databases through replication and sharding. You've seen how sharding in both relational databases and MongoDB allows you to distribute data across multiple machines, but we ended on a critical question: how do you maintain data consistency when a single operation, like placing an order, needs to update data on several different shards or services?
This lesson directly addresses that challenge. We'll explore the complexities of distributed transactions, where atomicity—the all-or-nothing guarantee—is much harder to achieve. Our focus will be on the Two-Phase Commit (2PC) protocol, a foundational algorithm for ensuring strong consistency in distributed systems. By the end of this lesson, you will be able to explain the core challenges of distributed transactions and the precise role that 2PC plays in solving them, including its significant trade-offs.
1. The Challenge of Distributed Atomicity
In a single database, transactions are a solved problem. The database engine ensures that a series of operations either all succeed (commit) or are all undone (rollback), maintaining the ACID properties. But what happens when your services are distributed, as is common in modern architectures?
Consider a food delivery service. Placing an order might involve:
- The
Order Servicecreating an order record. - The
Store Servicereserving the food item from inventory. - The
Delivery Serviceassigning a delivery partner.
If these are three separate services, each with its own database, what happens if the Store Service successfully reserves the food, but the Delivery Service fails to find a partner? You're left in an inconsistent state. The food is reserved but can't be delivered, leading to waste and a poor user experience. This is the core problem of distributed transactions: ensuring atomicity across multiple, independent systems.
The following video uses this exact scenario to explain the problem in a very clear and practical way.
Distributed Transactions: Two-Phase Commit Protocol
This video, "Distributed Transactions: Two-Phase Commit Protocol" by Arpit Bhayani, provides an excellent real-world analogy using Zomato's food delivery business.
Watch the initial segments that set up the problem, from the introduction through the explanation of the engineering challenges. Focus on why simple sequential HTTP calls between services are not enough to guarantee an atomic outcome.
As the video highlights, if any step in the process fails, the entire operation must be aborted, and any completed steps must be undone. This "all-or-nothing" behavior across distributed components is what protocols like Two-Phase Commit aim to provide.
2. The Two-Phase Commit (2PC) Protocol
Two-Phase Commit is a classic protocol for achieving atomic commitment in a distributed system. It introduces a central Coordinator that orchestrates the transaction across multiple Participants (e.g., database shards, microservices).
The protocol, as its name suggests, operates in two distinct phases:
-
Phase 1: The Prepare Phase (Voting)
- The Coordinator sends a
preparemessage to all Participants, asking them if they are ready to commit the transaction. - Each Participant attempts to perform its part of the transaction. It secures all necessary resources (e.g., by acquiring locks) and writes the changes to a temporary, durable log.
- If a Participant is confident it can complete the transaction, it responds with a
YESvote. This is a promise: once a Participant votesYES, it cannot unilaterally back out. It must wait for the Coordinator's final decision. - If a Participant cannot complete its work, it responds with a
NOvote.
- The Coordinator sends a
-
Phase 2: The Commit Phase (Decision)
- Success Case: If the Coordinator receives
YESvotes from all Participants, it knows the transaction can succeed. It records aCOMMITdecision in its own durable log and then sends acommitmessage to all Participants. Participants then make their changes permanent and release any locks. - Failure Case: If the Coordinator receives even one
NOvote, or if a Participant times out, it knows the transaction must fail. It records anABORTdecision and sends arollbackorabortmessage to all Participants. Participants then use their logs to undo any changes they prepared.
- Success Case: If the Coordinator receives
The following diagram illustrates the flow for both the Coordinator and a Participant.

The following article provides a clear textual explanation of these phases.
Demystifying Two Phase Commit (2PC) for Distributed ...
This article from DEV Community clearly breaks down the two phases of the 2PC protocol.
Please read the sections What is 2PC and the subsequent Decision Phase. The diagrams for the commit and rollback scenarios are particularly helpful for visualization.
3. The Core Challenge of 2PC: It's a Blocking Protocol
While 2PC provides the strong consistency (atomicity) we need, it comes with significant drawbacks that are critical to understand for system design interviews. The most severe limitation is that 2PC is a blocking protocol.
Consider this scenario:
- The Coordinator sends
preparerequests to all Participants. - All Participants successfully prepare, write to their logs, acquire locks, and vote
YES. - The Coordinator crashes before it can send the final
commitorabortdecision to the Participants.
In this state, the Participants are blocked. They have promised to commit if asked, so they cannot unilaterally abort. But they also cannot commit, because they don't know the final decision. They must wait for the Coordinator to recover. While blocked, they continue to hold locks on database rows or other resources, which can bring parts of your system to a halt. If the coordinator is down for a long time, these resources remain locked, impacting availability.
Martin Kleppmann, a leading author in distributed systems, explains this problem eloquently.
Distributed Systems 7.1: Two-phase commit
This excerpt from Martin Kleppmann's distributed systems lecture series provides a precise, academic explanation of 2PC and its primary weakness.
First, watch the section explaining the mechanics of 2PC. Then, pay close attention to the following segment, which discusses what happens if the coordinator crashes. This highlights the core "blocking" problem.
This blocking nature, combined with the coordinator being a single point of failure and a potential performance bottleneck, makes standard 2PC a challenging pattern to use in high-throughput, highly available microservice architectures.
4. Overcoming the Blocking Problem (Advanced)
For a senior engineer, it's valuable to know that variants of 2PC exist to overcome the blocking problem. As Kleppmann briefly introduces, one such approach involves replacing the centralized coordinator with a fault-tolerant consensus algorithm like Raft or Paxos, using a total order broadcast.
In this model:
- Instead of sending votes to a single coordinator, Participants broadcast their
YES/NOvotes to all other Participants using the consensus system. - All Participants receive the same votes in the same order.
- They can then independently and consistently arrive at the same decision: commit if all votes are
YES, abort otherwise. - A failure detector is used to vote
NOon behalf of any crashed Participants.
This decentralized approach removes the single point of failure and the blocking problem, but at the cost of increased complexity. You can explore this in more detail by re-watching the final part of Kleppmann's video.
Distributed Systems 7.1: Two-phase commit
This final segment of the video explains how a non-blocking variant of 2PC can be implemented using total order broadcast.
Watch from the start of this section to understand how a consensus algorithm can replace the central coordinator to build a fault-tolerant atomic commit protocol.
5. 2PC in Practice: PostgreSQL and System Design Trade-offs
Despite its limitations, 2PC is implemented in many traditional relational databases. Since you have experience with PostgreSQL, it's useful to see how it's supported there. PostgreSQL provides specific SQL commands to manage the two phases.
Demystifying Two Phase Commit (2PC) for Distributed ...
The same DEV Community article has a great section on using 2PC in PostgreSQL.
Read the section on PostgreSQL syntax. Note the PREPARE TRANSACTION, COMMIT PREPARED, and ROLLBACK PREPARED commands. This shows how the two phases are exposed directly in SQL.
Given these trade-offs, when should a system designer choose 2PC? The answer lies in the required level of consistency. 2PC is often compared to the Saga pattern, an alternative for managing long-running transactions in microservices.

This table perfectly summarizes the choice:
- Choose 2PC when you need strong (ACID) consistency and the transaction is relatively short-lived. It's suitable for critical operations where temporary inconsistencies are unacceptable (e.g., a core financial transfer between two ledgers).
- Choose Saga when you can tolerate eventual consistency and need high availability and resilience. The Saga pattern works by executing a sequence of local transactions, where each step has a corresponding compensating action to undo it if a later step fails. This is non-blocking and a better fit for most long-running business processes in a microservices world.
This article provides a final, concise guide on making this crucial design decision.
Distributed Transactions (2PC, Saga) in System Design - DEV Community
This article directly addresses the decision-making process for a system designer.
First, read the section summarizing the Limitations of 2PC. Then, read the final section, Choosing Between 2PC and Saga, which frames the decision in the context of microservices and system design trade-offs.
Conclusion
In this lesson, we dissected the challenge of maintaining atomicity in distributed systems. You've learned how the Two-Phase Commit protocol provides a solution for strong consistency, but at a significant cost to availability and performance.
Key Takeaways:
- Distributed Atomicity is Hard: Simple sequential operations across services can lead to inconsistent states if one part fails.
- 2PC Provides Strong Consistency: The Coordinator-led prepare and commit phases ensure that a distributed transaction is atomic (all-or-nothing).
- 2PC's Main Drawback is Blocking: If the Coordinator fails after the prepare phase, Participants are left in a blocked state, holding resources and waiting for recovery. This makes the Coordinator a single point of failure and impacts system availability.
- 2PC vs. Saga is a Key Design Choice: For system design interviews, you should be able to articulate the trade-off. 2PC is for short-lived, critical transactions needing immediate consistency. The Saga pattern is generally preferred for long-running business processes in microservices that prioritize availability and can tolerate eventual consistency.
We've focused today on the complexities of guaranteeing consistency for writes. However, a huge part of building scalable systems is optimizing for reads. In our next lesson, we will shift gears to this topic and explore how to drastically reduce latency and database load by introducing a distributed cache.