8. Transactions
Chapter 8 of Designing Data-Intensive Applications — 13 sections.
"Some authors have claimed that general two-phase commit is too expensive to support… We believe it is better to have application programmers deal with performance problems due to overuse of transactions as bottlenecks arise, rather than always coding around the lack of transactions." — James Corbett et al., Spanner: Google's Globally-Distributed Database (2012)
The six things that go wrong, which transactions exist to hide:
- The database software or hardware may fail at any time — including in the middle of a write
- The application may crash at any time — including halfway through a series of operations
- Network interruptions can cut off the application from the database, or one node from another
- Several clients may write at the same time, overwriting one another's changes
- A client may read data that doesn't make sense because it has only partially been updated
- Race conditions between clients can cause surprising bugs
A transaction groups several reads and writes into a logical unit. Conceptually all of them execute as ONE operation: either the entire transaction succeeds (COMMIT) or it fails (ABORT / ROLLBACK). If it fails, the application can safely RETRY.
Transactions are not a law of nature; they were created with a purpose — to simplify the programming model for applications accessing a database. Using them lets the application ignore certain error scenarios and concurrency issues because the database handles them instead. We call these safety guarantees.
And the counterpoint, stated fairly: not every application needs transactions, and sometimes there are advantages to weakening or abandoning them (better performance, higher availability). Some safety properties can be achieved without transactions. On the other hand — the technical cause behind the Post Office Horizon scandal (Ch 2) was probably a lack of ACID transactions in the underlying accounting system.
The historical arc, corrected
| Era | What happened to transactions |
|---|---|
| 1975 | IBM System R — the first SQL database — defines the transaction model. “The general idea has remained virtually the same for 50 years: transaction support in MySQL, PostgreSQL, Oracle, SQL Server is uncannily similar to System R.” |
| Late 2000s | NoSQL rises, with new data models and replication and sharding by default. “Transactions were the main casualty: many databases abandoned them entirely, or redefined the word to describe a much weaker set of guarantees.” The hype produced a popular belief that transactions were fundamentally unscalable. |
| Now | “More recently, that belief has turned out to be wrong.” NewSQL — CockroachDB, TiDB, Spanner, FoundationDB, YugabyteDB — show that transactional systems scale to large data volumes and high throughput, by combining sharding with consensus protocols (Ch 10). |
Sections
- 8.12The Meaning of ACIDCoined 1983 by Theo Härder and Andreas Reuter, to establish precise terminology for fault-tolerance mechanisms.
- 8.21Single-Object and Multi-Object OperationsToo slow with many emails → denormalize into a separate counter field, incremented on new mail and decremented on read.
- 8.36Weak Isolation Levels
- 8.47Serializability
- 8.56Distributed Transactions
- 8.6THE ANOMALY TABLE — the single most useful table in the book
- 8.7Deep divesTechnology deep dives
- 8.8Failure catalogProduction failure catalog for this chapter
- 8.9Decision sheet1Decision cheat sheetPrefer, in order: (1) redesign so the transaction fits in one shard; (2) idempotency + a message-ID table (§5.6) — this covers most "exactly-once" needs with only local transactio…
- 8.10Worked examples1Worked examples
- 8.11Self-testSelf-test
- 8.12TerminologyTerminology introduced here
- 8.13Forward linksForward links