Learn Labs
8. Transactions

8.4 Serializability

The bleak summary the chapter arrives at:

  • Isolation levels are hard to understand and inconsistently implemented (the meaning of "repeatable read" varies significantly)
  • It can be difficult to tell BY LOOKING AT THE APPLICATION CODE whether it is safe to run at a particular isolation level — especially in a large application where you might not know everything happening concurrently
  • There are NO GOOD TOOLS to help us detect race conditions. Static analysis may help in principle, but research techniques have not yet found their way into practical use. Testing is hard, because problems occur only if you get unlucky with the timing

This has been the situation since the 1970s. All along, the answer from researchers has been simple: USE SERIALIZABLE ISOLATION.

Three implementations:

Actual serial executionTwo-phase locking (2PL)Serializable snapshot isolation (SSI)
Pessimistic to the extremePessimisticOptimistic
One transaction at a time, on a single thread.Readers block writers, writers block readers.No blocking; check at commit time and abort if the execution wasn’t serializable.
VoltDB/H-Store, Redis, DatomicMySQL/InnoDB and SQL Server serializable; Db2 repeatable read. “For around 30 years, the only algorithm widely used.”PostgreSQL serializable, SQL Server Hekaton, HyPer, CockroachDB, FoundationDB, BadgerDB. First described in 2008.

4.1 Actual Serial Execution

Two developments made it viable, only in the 2000s:

  1. RAM became cheap enough to keep the entire active dataset in memory — when all data is in memory, transactions execute much faster
  2. Designers realized OLTP transactions are usually SHORT and make only a SMALL number of reads and writes. Long-running analytical queries are typically read-only, so they can run on a consistent snapshot OUTSIDE the serial execution loop

A system designed for single-threaded execution can sometimes PERFORM BETTER than one supporting concurrency, because it avoids the coordination overhead of locking. However, its throughput is LIMITED TO THAT OF A SINGLE CPU CORE.

Why stored procedures are mandatory here

The history: designers originally intended a transaction to encompass an entire flow of user activity — booking an airline ticket: search routes/fares/seats, decide, book seats on each flight, enter passenger details, pay. Unfortunately, HUMANS ARE VERY SLOW TO MAKE UP THEIR MINDS. A transaction waiting for user input would require a huge number of concurrent, mostly idle transactions. So almost all OLTP applications keep transactions short — on the web, a transaction is committed within the same HTTP request; a new request starts a new transaction.

But even with humans removed, the interactive client/server style remains:

Interactive transaction — a round trip per statement
appdatabaseBEGINround tripquery 1result — round tripquery 2result — round tripCOMMIT

“A lot of time is spent in network communication.” With no concurrency, throughput would be dreadful, because the database spends most of its time waiting for the application to issue the next query.

Stored procedure — one call
appdatabasecall procedureresult

The database executes the entire transaction with no waiting for the network or for disk I/O — very fast, provided all the required data is in memory.

Figure 8.4.2But even with humans removed, the interactive client/server style remains

Therefore systems with single-threaded serial transaction processing DON'T ALLOW interactive multistatement transactions. The application must either limit itself to single-statement transactions or submit the entire transaction code ahead of time as a STORED PROCEDURE.

Stored procedures' bad reputation — four reasons, all fair:

  1. Each vendor had its own language (PL/SQL, T-SQL, PL/pgSQL) which haven't kept up with general-purpose languages — ugly and archaic, lacking the ecosystem of libraries
  2. Code running in a database is difficult to manage: harder to debug, awkward to version control and deploy, trickier to test, difficult to integrate with a metrics collection system
  3. A database is much more performance-sensitive than an application server, because one instance is shared by many app servers. A badly written stored procedure can cause much more trouble than equivalent bad code in an app server
  4. In a multitenant system allowing tenants to write stored procedures, it's a SECURITY RISK to execute untrusted code in the same process as the database kernel

But modern implementations fixed #1: VoltDB uses Java or Groovy, Datomic uses Java or Clojure, Redis uses Lua, MongoDB uses JavaScript.

(A modern use case the book calls out: GraphQL proxies exposing the database directly. If the proxy doesn't support complex validation logic, you can embed it in a stored procedure — otherwise you must deploy a validation service between the proxy and the database.)

VoltDB also uses stored procedures for REPLICATION: instead of copying writes, it executes the same stored procedure on each replica. This REQUIRES stored procedures to be DETERMINISTIC — a transaction needing the current date/time must use special deterministic APIs. This is STATE MACHINE REPLICATION (Ch 10).

Sharding for serial execution

To scale beyond one core, shard. If each transaction reads and writes within a SINGLE SHARD, each shard gets its own transaction processing thread — give each CPU core its own shard and throughput scales LINEARLY with cores.

But any transaction touching multiple shards must be coordinated across all of them — the stored procedure must be performed IN LOCKSTEP across all shards.

VoltDB reports about 1,000 CROSS-SHARD writes per second — ORDERS OF MAGNITUDE below its single-shard throughput, AND IT CANNOT BE INCREASED BY ADDING MORE MACHINES.

Whether transactions can be single-shard depends on the data: simple key-value data shards easily; data with multiple secondary indexes is likely to require a lot of cross-shard coordination (Ch 7 §5).

The four constraints, summarized:

  1. Every transaction must be small and fast — IT TAKES ONLY ONE SLOW TRANSACTION TO STALL ALL TRANSACTION PROCESSING
  2. The active dataset must fit in memory. Rarely accessed data could go to disk, but if a single-threaded transaction needed it, the system would get very slow
  3. Write throughput must fit on one CPU core, or transactions must shard without cross-shard coordination
  4. Cross-shard transactions are possible, but their throughput is hard to scale

4.2 Two-Phase Locking (2PL)

⚠️ 2PL IS NOT 2PC. 2PL provides serializable ISOLATION; 2PC provides atomic COMMIT in a distributed database. Best to think of them as entirely separate concepts and ignore the unfortunate similarity in the names.

The rule:

Several transactions may concurrently read the same object — as long as nobody is writing to it. But as soon as anyone wants to write, exclusive access is required.

SituationConsequence
A has read an object, B wants to write itB waits for A to commit or abort, so B can’t change it behind A’s back.
A has written an object, B wants to read itB waits for A to commit or abort — reading an old version is not acceptable under 2PL.

“Writers don’t just block other writers; they also block readers, and vice versa.” That is the exact opposite of snapshot isolation’s “readers never block writers”.

Implementation — a shared/exclusive (multi-reader single-writer) lock on each object:

  • Read → acquire the lock in shared mode. Several transactions may hold it in shared mode; if another holds it exclusively, wait
  • Write → acquire the lock in exclusive mode. No other transaction may hold it at all
  • Read then write → upgrade the shared lock to exclusive
  • Hold every lock until the END of the transaction (commit or abort)

This is where "two-phase" comes from: the first phase (GROWING) is when locks are acquired while the transaction executes; the second phase (SHRINKING) is when all locks are released at the end. THE TWO PHASES MUST NOT OVERLAP — once a lock is released, no new locks may be acquired.

Deadlock is frequent. The database automatically detects deadlocks and aborts one transaction; the application must retry.

Performance — why it hasn't been the default since the 1970s

Transaction throughput and response times are SIGNIFICANTLY WORSE under 2PL. Partly the overhead of acquiring and releasing locks — but MORE IMPORTANTLY, REDUCED CONCURRENCY. By design, if two concurrent transactions try to do anything that MAY IN ANY WAY result in a race condition, one has to wait.

The pathological case:

  1. A transaction that reads an entire table — a backup, an analytical query, an integrity check — must take a shared lock on the entire table.
  2. It first waits until all in-progress writers to that table complete.
  3. Then, while the whole table is read (which may take a long time), all other transactions wanting to write are blocked.

“In effect, the database becomes unavailable for writes for an extended time.”

Databases running 2PL can have quite UNSTABLE LATENCIES, and can be VERY SLOW AT HIGH PERCENTILES if there is contention. JUST ONE SLOW TRANSACTION, or one that accesses a lot of data and acquires many locks, could cause the rest of the system to grind to a halt. Transaction timeouts and slow-query monitoring are used to detect and limit misbehaving queries.

And deadlocks occur MUCH more frequently under 2PL than under lock-based read-committed. When a deadlocked transaction is aborted and retried, IT NEEDS TO DO ITS WORK ALL OVER AGAIN — if deadlocks are frequent, significant wasted effort.

Predicate locks — how 2PL handles phantoms

A database with serializable isolation MUST prevent phantoms. Conceptually you need a PREDICATE LOCK — like a shared/exclusive lock, but rather than belonging to a particular object, it belongs to ALL OBJECTS THAT MATCH A SEARCH CONDITION.

SELECT * FROM bookings
  WHERE room_id = 123 AND end_time > '2026-01-01 12:00'
                      AND start_time < '2026-01-01 13:00';
  • A wants to READ objects matching a condition → acquire a shared-mode predicate lock on the query's conditions. If B holds an exclusive lock on any matching object, A waits
  • A wants to INSERT/UPDATE/DELETE → first check whether EITHER THE OLD OR THE NEW VALUE matches any existing predicate lock. If B holds a matching one, A waits

The key idea: A PREDICATE LOCK APPLIES EVEN TO OBJECTS THAT DO NOT YET EXIST IN THE DATABASE but might be added in the future (phantoms). If 2PL includes predicate locks, the database prevents all forms of write skew and other race conditions — its isolation becomes serializable.

Index-range locks (next-key locking) — what's actually implemented

Predicate locks DO NOT PERFORM WELL: with many locks held by active transactions, CHECKING FOR MATCHING LOCKS BECOMES TIME-CONSUMING. So most 2PL databases implement index-range locking, a simplified APPROXIMATION.

It's SAFE to simplify a predicate by making it match a GREATER set of objects. A lock on bookings of room 123 between noon and 1 pm can be approximated by locking bookings for room 123 at ANY time, or locking ALL ROOMS between noon and 1 pm. Any write matching the original predicate will definitely also match the approximation.

IndexWhat gets locked
Index on room_idA shared lock on the index entry for 123 — “a transaction has searched for bookings of room 123”.
Index on start/end timeA shared lock on a range of values in that index — “a transaction has searched for bookings overlapping noon–1 pm on that date”.

Another transaction inserting, updating or deleting a booking for the same room and/or an overlapping period must update the same part of the index, encounters the shared lock, and is forced to wait.

Index-range locks are NOT AS PRECISE as predicate locks (they may lock a bigger range than strictly necessary), but since they have MUCH LOWER OVERHEADS, they are a GOOD COMPROMISE.

If there is NO SUITABLE INDEX to attach a range lock to, the database FALLS BACK TO A SHARED LOCK ON THE ENTIRE TABLE. Not good for performance — it stops all other writes — but a SAFE fallback.

(Practical corollary: a missing index doesn't just make a serializable query slow — it makes it lock the whole table.)

4.3 Serializable Snapshot Isolation (SSI)

Are serializable isolation and good performance fundamentally at odds? IT SEEMS NOT: SSI provides FULL SERIALIZABILITY with only a SMALL PERFORMANCE PENALTY compared to snapshot isolation. First described 2008.

Pessimistic vs optimistic
Principle
2PL — pessimisticIf anything MIGHT possibly go wrong (indicated by a lock held by another transaction), it's better to WAIT until the situation is safe. Like mutual exclusion in multithreaded programming
Serial execution — pessimistic to the extremeEssentially equivalent to each transaction holding an exclusive lock on the ENTIRE DATABASE (or shard). We compensate by making each transaction very fast, so it holds the "lock" only briefly
SSI — optimisticInstead of blocking when something potentially dangerous happens, transactions CONTINUE ANYWAY, in the hope that everything will turn out all right. At commit, the database checks whether isolation was violated; if so, ABORT AND RETRY. ONLY TRANSACTIONS THAT EXECUTED SERIALIZABLY ARE ALLOWED TO COMMIT

When optimistic loses: it performs badly under HIGH CONTENTION (many transactions accessing the same objects), leading to a high proportion of aborts. If the system is already close to maximum throughput, THE ADDITIONAL LOAD FROM RETRIED TRANSACTIONS CAN MAKE PERFORMANCE WORSE.

When it wins: with enough spare capacity and not-too-high contention, optimistic tends to perform BETTER than pessimistic.

Reducing contention: commutative atomic operations — several transactions incrementing a counter don't care about order (as long as the counter isn't read in the same transaction), so the concurrent increments can all be applied without conflicting.

The core idea: decisions based on an outdated premise

The recurring write-skew pattern: a transaction reads data, examines the result, and DECIDES to take an action based on what it saw. Under snapshot isolation, THE RESULT MAY NO LONGER BE UP TO DATE BY THE TIME THE TRANSACTION COMMITS.

The transaction is acting on a PREMISE — "there are currently two doctors on call." Later, the premise may no longer be true.

The database doesn't know how the application logic uses the query result. To be safe, IT MUST ASSUME THAT ANY CHANGE IN THE QUERY RESULT MEANS WRITES IN THAT TRANSACTION MAY BE INVALID — there may be a causal dependency between the queries and the writes. So it must detect situations where a transaction MAY HAVE ACTED ON AN OUTDATED PREMISE and abort.

Two detection cases:

(a) Detecting stale MVCC reads — the write happened BEFORE the read, but committed after
txn 42databasetxn 43BEGINAaliyah.on_call = false (uncommitted)SELECT on_callAaliyah is on callCOMMIT — 42 commits firstBryce.on_call = falseCOMMIT → ABORTED

Txn 43 ignores 42’s uncommitted write, as the MVCC visibility rules require. But once 42 commits, that ignored write has now taken effect, so 43’s premise is no longer true — and 43 is aborted at commit time.

Figure 8.4.6Two detection cases

The database TRACKS when a transaction ignores another transaction's writes because of MVCC visibility rules. At commit, it checks whether any ignored writes have NOW been committed. If so, abort.

Why wait until commit rather than aborting immediately? Three reasons:

  1. If transaction 43 were READ-ONLY, it wouldn't need to abort — there's no risk of write skew. At read time the database doesn't yet know whether it will later write
  2. Transaction 42 may yet ABORT, or may still be uncommitted when 43 commits — so the read may turn out NOT to have been stale after all
  3. By avoiding unnecessary aborts, SSI PRESERVES SNAPSHOT ISOLATION'S SUPPORT FOR LONG-RUNNING READS from a consistent snapshot
(b) Detecting writes that affect prior reads — the write happens AFTER the read

Both 42 and 43 search for on-call doctors during shift 1234. If there is an index on shift_id, the database uses index entry 1234 to record the fact that 42 and 43 read this data. With no index, it tracks reads at the table level.

txn 42index entry 1234txn 43SELECT … shift_id = 1234SELECT … shift_id = 1234Aaliyah.on_call = falseBryce.on_call = falseyour read may be staleCOMMITCOMMIT → ABORT

When writing, the database looks in the indexes for other transactions that recently read this data. It is like acquiring a write lock — but rather than blocking, the lock acts as a tripwire: it notifies those transactions that the data they read may no longer be up to date.

42 commits successfully because 43 hasn’t committed yet, so 43’s write hasn’t taken effect. 43 aborts, because 42’s conflicting write has already committed.

Figure 8.4.7(b) Detecting writes that affect prior reads — the write happens AFTER the read

This information is kept only for a while: after a transaction finishes and all concurrent transactions finish, the database can forget what it read.

Performance of SSI

Granularity trade-off: detailed tracking → precise aborts but significant bookkeeping overhead. Less detailed → faster, but more transactions aborted than strictly necessary.

PostgreSQL reduces unnecessary aborts using theory showing it's sometimes OK for a transaction to read information overwritten by another — depending on what else happened, the execution may still be provably serializable.

Compared toSSI's advantage
2PLOne transaction doesn't block waiting for locks held by another. Writers don't block readers, and vice versa. This makes QUERY LATENCY MUCH MORE PREDICTABLE AND LESS VARIABLE. Read-only queries can run on a consistent snapshot WITHOUT ANY LOCKS — very appealing for read-heavy workloads
Serial executionNOT limited to a single CPU core. FoundationDB DISTRIBUTES the detection of serialization conflicts across multiple machines, scaling to very high throughput; transactions can read and write in multiple shards while ensuring serializable isolation
Nonserializable SISome overhead. How significant is a matter of debate: some believe serializability checking is not worth it; others believe its performance is now so good that there is no need to use the weaker snapshot isolation anymore

The RATE OF ABORTS significantly affects overall performance. A transaction that reads and writes over a LONG PERIOD is likely to conflict and abort — SSI requires READ/WRITE transactions to be fairly SHORT (long-running READ-ONLY transactions are fine). However, SSI is LESS SENSITIVE TO SLOW TRANSACTIONS than 2PL or serial execution.


On this page