Transactions and Concurrency Control

Topics Covered

The Purpose of Transactions

Atomicity: All or Nothing

Consistency: Your Rules, Not the Database's

Isolation: Pretending You Are Alone

Durability: Surviving the Worst

Isolation Levels

Read Committed: The Baseline

Snapshot Isolation and MVCC

What Snapshot Isolation Does Not Prevent

Implementing Serializable Isolation

Actual Serial Execution

Two-Phase Locking (2PL)

Serializable Snapshot Isolation (SSI)

Distributed Transactions

Two-Phase Commit (2PC)

The Coordinator as Single Point of Failure

XA Transactions

Beyond Transactions

Exactly-Once vs Effectively-Once Processing

When Transactions Are Too Expensive

Databases fail. Disks corrupt. Networks drop. Processes crash mid-operation. Transactions exist because real systems are hostile environments where anything can go wrong at any time. A transaction groups multiple reads and writes into a single logical unit: either all of them succeed, or none of them do. Without this guarantee, a crash halfway through writing an order could leave your database with a debited account but no corresponding order record, money vanishing into the void.

The transaction concept was born in the 1970s when database systems first needed to handle multiple concurrent users modifying shared data. Before transactions, applications had to implement their own crash recovery and concurrency control, a task so error-prone that the same bugs were reinvented at every company. Transactions moved this complexity into the database, giving application developers a simple mental model: wrap your operations in BEGIN and COMMIT, and the database handles the rest.

This simplification is one of the most important abstractions in computer science. A developer writing application code does not need to think about disk failures, concurrent access, or partial writes. The transaction boundary handles all of that. The application says "these operations must happen together," and the database figures out how to make that true despite hardware failures, concurrent access, and everything else that can go wrong.

The properties that define a transaction are captured by the acronym ACID. Understanding what each letter actually means (and does not mean) is essential because databases implement these properties at very different levels of strictness. The term "ACID" was coined by Theo Haerder and Andreas Reuter in 1983, but the individual properties had been implemented in systems like IBM's System R years earlier.

The same interleaving run at read uncommitted and read committed, with the uncommitted value only one of them exposes.

Atomicity: All or Nothing

Atomicity does not mean "happening at one instant." It means abortability. If a transaction writes three rows and crashes after the second write, atomicity guarantees the first two writes are rolled back. The database returns to the state before the transaction began. The client can safely retry the entire operation knowing that the partial failure left no trace.

This is implemented through a write-ahead log (WAL). Before modifying any data page, the database writes the intended change to a sequential log on disk. On crash recovery, the database replays committed transactions from the log and undoes uncommitted ones. The WAL is what makes "undo" possible.

Consider a concrete example. An e-commerce checkout transaction does four things: inserts an order row, inserts order-item rows, decrements inventory, and charges the payment. If the process crashes after decrementing inventory but before inserting the order, atomicity rolls back the inventory change. Without atomicity, you would need application-level cleanup code to detect and repair every possible partial failure state, a combinatorial nightmare that grows with each new step in the transaction.

The terminology can be confusing. In the context of concurrent programming, "atomic" often means "executes as a single indivisible operation" (like an atomic CPU instruction). In the context of databases, "atomic" means "abortable with complete rollback." A database transaction is not indivisible; it consists of many individual operations spread over time. What atomicity guarantees is that the effects of those operations are indivisible: observers see either all of them or none of them, never a subset.

Consistency: Your Rules, Not the Database's

Consistency is the odd one out in ACID. It is not a database property. It is an application property. Consistency means your data satisfies whatever invariants your application requires: account balances must not go negative, every order must reference a valid customer, foreign keys must point to existing rows.

The database provides tools (constraints, foreign keys, triggers) to enforce some invariants, but ultimately consistency depends on the application writing correct transaction logic. A database cannot know that your business rule prohibits more than 5 items per cart. Only your application code enforces that.

This distinction matters because developers sometimes believe that using a database with "ACID compliance" automatically prevents all data corruption. It does not. The database guarantees that your transaction runs atomically, in isolation, and durably. But if your transaction logic is wrong (e.g., it debits without checking the balance first), the database faithfully executes the incorrect logic. Consistency is a shared responsibility between the database engine and your code.

Some databases blur this line by offering features that enforce more invariants automatically. PostgreSQL's CHECK constraints can enforce "balance >= 0" at the database level. Foreign key constraints ensure referential integrity. Unique indexes prevent duplicate entries. The more invariants you push into the database schema, the harder it is for buggy application code to violate them. But complex business rules (monthly spending limits, multi-entity invariants, temporal constraints) almost always require application-level enforcement.

Isolation: Pretending You Are Alone

Isolation means concurrent transactions do not interfere with each other. The ideal is that every transaction executes as if it were the only transaction running on the database. In practice, this ideal (called serializability) is expensive, so databases offer weaker isolation levels that trade correctness for performance. Most bugs caused by concurrency happen because developers assume stronger isolation than their database actually provides.

Imagine two users simultaneously buying the last item in stock. Without isolation, both transactions read "1 item remaining," both decrement to 0, and both succeed, selling an item you do not have. With proper isolation, one transaction sees the item, buys it, and commits. The second transaction either sees "0 items remaining" and fails gracefully, or is blocked until the first completes. The isolation level determines which of these behaviors your database actually provides.

Common Pitfall

The most dangerous misconception in database engineering is assuming your database provides serializable isolation by default. Most databases default to read committed or repeatable read. This means concurrent transactions CAN interfere with each other in subtle ways. Know your database's default isolation level and what anomalies it permits.

Durability: Surviving the Worst

Durability means that once a transaction commits, its data survives any subsequent crash, power failure, or hardware fault. This sounds simple, but the implementation involves flushing the WAL to disk (fsync), and in replicated databases, waiting for the write to reach multiple nodes before acknowledging the commit.

On a single node, durability means the data is on a persistent storage device. The database calls fsync to force the operating system to flush buffered writes to the physical disk. Without fsync, a power failure could lose data that the application believed was safely written because it was sitting in the OS page cache, not on the actual disk platters or flash cells. Some databases offer configurable durability: PostgreSQL's synchronous_commit = off setting acknowledges the commit before fsync completes, reducing latency at the risk of losing the last few milliseconds of committed transactions on crash. This tradeoff is acceptable for some workloads (logging, session data) but dangerous for others (financial records).

In a replicated database, durability extends beyond a single machine. A transaction is not considered committed until the write has been replicated to a configurable number of nodes (the write quorum). If the primary node catches fire, a replica can take over with the committed data intact.

The choice of write quorum determines both durability and latency. A quorum of 1 (write to primary only, replicate asynchronously) gives the lowest latency but risks data loss if the primary fails before replication. A quorum of a majority (e.g., 2 out of 3 nodes) provides strong durability but adds the latency of waiting for the slowest node in the quorum. Neither single-node nor multi-node durability is absolute: if every disk in every replica fails simultaneously, data is lost. Durability is a spectrum, not a binary.

Two different things are called a quorum across this course and they are worth separating now. A replication quorum is tunable per request and is about how many copies exist, with the condition W+R>NW + R > N ensuring a reader overlaps a writer; it is developed in replication strategies. A consensus quorum is always a majority and is about which of several competing proposals is the decision; it is developed in consistency and consensus. The first gives you durability and read-your-writes. Only the second gives you agreement on order, which is why a Dynamo-style quorum with W+R>NW + R > N still produces concurrent conflicting versions and a Raft cluster does not.

Interview Tip

In interviews, when explaining ACID, focus on atomicity and isolation. Consistency is an application concern, and durability is handled by the storage engine. The interesting design decisions live in how atomicity handles partial failures (WAL, undo logs) and how isolation balances correctness against performance (isolation levels).