skip to content

Google's Percolator system implements multi-row, multi-table ACID transactions on top of Bigtable, a key-value store with no native multi-row transaction support, without using a classic centralized two-phase-commit coordinator process. Describe how Percolator achieves this using per-row locks and timestamps, and what it trades away compared to a traditional 2PC coordinator.

level: principalimportance: nice to knowfreq 20%

answer

  1. lock+write columns on rows
  2. primary row = atomic commit point
  3. timestamp oracle for ordering
  4. secondaries cleaned up lazily/async
  5. snapshot isolation, not serializable

basics

~20 s

Percolator gives every transaction a start time and commit time, and instead of one central coordinator, it stores tiny 'lock' markers directly in the data rows themselves, with one row acting as the primary lock that decides the whole transaction's fate. Other transactions that bump into a lock can check on the primary row to figure out what happened.

solid answer

~60 s

Percolator layers snapshot-isolation transactions on Bigtable by adding lock and write columns to every row and using a centralized timestamp oracle to hand out globally-ordered start and commit timestamps. A transaction reads at its start timestamp (seeing only committed data before that point), and to commit, it designates one of its written rows as the 'primary' and writes a lock pointing to the primary on every row, in parallel, rather than going through a single coordinator process. Once all locks are acquired (2PC's prepare phase, but distributed via the rows themselves), it commits the primary first, which is the atomic decision point, then asynchronously converts the remaining locks into committed writes. Because the primary's lock/commit record lives durably in Bigtable rather than a coordinator's private log, other transactions that encounter a stale lock can inspect the primary row directly to determine whether to roll it forward or clean it up, instead of blocking on a dead coordinator. The trade is higher per-transaction latency (extra Bigtable round trips replace a lightweight coordinator conversation) and only snapshot-isolation-level guarantees, in exchange for removing the single-coordinator bottleneck and failure mode, and for horizontal scalability.

go deeper

for a junior

Not expected to know Percolator by name; at most should recognize it's a way to do transactions across multiple rows in a key-value store.

for a middle

Should grasp the basic idea that Percolator stores lock information in the data itself instead of a separate coordinator, if prompted.

for a senior

Should be able to describe the primary/secondary row mechanism and why encountering a stale lock doesn't require contacting a live coordinator.

for a principal

Should explain the full mechanism (timestamp oracle, prewrite/commit phases, lazy secondary cleanup), articulate the latency-vs-availability trade against classic 2PC, and connect it to related systems like Spanner.

## What Percolator is for **Percolator** is a real, published system (Google's 2010 paper 'Large-scale Incremental Processing Using Distributed Transactions and Notifications') built to let Google's web search indexing pipeline make small, incremental, cross-row updates to its multi-petabyte index — stored in Bigtable — with the same atomicity guarantees a relational database gives you, but at a scale and update frequency no single relational database or classic 2PC coordinator could sustain. Its central design idea is: don't build a separate coordinator process at all — embed the coordination state directly into the rows of the data store itself, and use timestamps to establish global ordering instead of locks that must be released by a central authority. ## The two extra pieces The mechanism relies on two pieces of extra infrastructure layered on top of plain Bigtable. 1. **First**, every logical row gets extra columns beyond its normal data: a 'lock' column (empty when unlocked, or pointing to the transaction's primary row when locked) and a 'write' column (recording, for each committed version, which data timestamp it corresponds to). 2. **Second**, a single small, highly-available service called the **timestamp oracle** hands out strictly increasing timestamps on request — nothing else about it is centralized in the transaction-processing sense; it doesn't track transaction state, just issues numbers, so it's cheap to keep highly available and doesn't become a coordination bottleneck. ## Reads, then Prewrite A transaction starts by requesting a start timestamp from the oracle and does all its reads as of that timestamp, giving it a consistent snapshot of everything committed before it started (snapshot isolation) — this is why concurrent transactions never block each other for reading. When it's ready to commit, it picks one of the rows it's writing to and designates it the 'primary'; every other written row is a 'secondary.' The prepare phase — Percolator calls this 'Prewrite' — writes a lock to every row (primary and secondaries) in parallel, each pointing back at the primary, and fails the whole transaction if any row it tries to lock already has a conflicting lock or has been written more recently than its start timestamp (a write-write conflict). This is functionally 2PC's prepare/vote phase, but there's no separate coordinator process collecting votes — the client driving the transaction does the prewrites directly, and 'success' is just 'all prewrites succeeded.' ## Where the commit actually lands The commit is where the atomicity actually lands: the client gets a commit timestamp from the oracle, and then commits the primary row first by writing its lock into a durable committed 'write' record at the commit timestamp and clearing the lock. That single atomic write to the primary row is the transaction's true, irrevocable commit point — analogous to the coordinator's durable decision record in classic 2PC, except it lives inside the data itself rather than in a separate coordinator's private log. After the primary is committed, the client goes back and converts each secondary's lock into a committed write record too, but this part can happen asynchronously and even be finished by someone else later, because the primary's state is now the single source of truth. ## Resolving a stale lock This is the elegant part: if any other transaction encounters a stale lock on a secondary row (because the original client crashed mid-commit before cleaning up the secondaries), it doesn't need to talk to a dead coordinator — it reads the primary row directly. - If the primary shows a committed write record, the encountering transaction rolls the stale secondary lock forward into a committed write itself (helping finish someone else's transaction). - If the primary shows no lock and no committed record, it rolls the lock back. Either way, resolution is self-service, using data already durably stored in Bigtable, rather than requiring a live coordinator process to be reachable — this is exactly the failure mode that makes classic 2PC block, and Percolator sidesteps it by making the 'coordinator's decision' just another durable row in the same store everyone already reads. ## What it trades away What Percolator trades away, compared to a lightweight, tightly-coupled 2PC coordinator conversation, is **latency** and **isolation strength**. - Every transaction now costs multiple extra Bigtable round trips (prewrite each row, commit the primary, then clean up each secondary) rather than a few short messages to an in-memory coordinator, so Percolator transactions are notably slower per-transaction than a traditional in-process 2PC — Google's own paper reports roughly a two-to-threefold latency increase for their indexing workload compared to their prior batch-based system, an acceptable trade given the huge gain in update latency (minutes instead of days) it enabled overall. - It also only provides snapshot isolation, not full serializability, so certain write-skew anomalies that a stricter isolation level would prevent are possible. In exchange, Percolator removes the single-coordinator bottleneck and blocking-on-a-dead-coordinator failure mode entirely, and scales horizontally with Bigtable itself — thousands of concurrent transactions across a cluster, with no central process anyone has to keep highly available beyond the lightweight timestamp oracle. This general pattern — durable, self-describing lock/intent records co-located with the data instead of a separate coordinator — was influential enough to show up again, refined, in Google's Spanner and in TiDB's transaction layer.

  • Why does Percolator only need the timestamp oracle to be highly available, not the transaction coordination logic itself?
    The oracle's only job is issuing strictly increasing numbers on request — it holds no per-transaction state and doesn't need to remember anything about who is doing what, so keeping it available is a much smaller and simpler problem than keeping a stateful transaction coordinator available. All the actual per-transaction coordination state (locks, primary pointers) lives durably in Bigtable rows themselves, spread across the whole cluster rather than concentrated in one process.
  • What happens if a client crashes after prewriting all rows but before committing the primary?
    No commit ever happened, since the primary's durable write record was never created, so any other transaction that later encounters one of the stale locks will check the primary, find no committed record, and safely roll the lock back — the transaction is treated as if it never committed, which is the correct outcome since it genuinely didn't.
  • Why is snapshot isolation an acceptable trade for Percolator's use case even though it's weaker than full serializability?
    Percolator was built for Google's search-indexing pipeline, where most transactions touch disjoint or loosely related rows (e.g., updating a document's index entries), so the specific write-skew anomalies that only serializability rules out are rare and low-impact for that workload. Given that, the extra coordination cost of enforcing full serializability wasn't worth paying for the incremental-indexing use case Percolator targeted.

Instead of hiring one referee to personally confirm every player is ready before blowing the whistle, Percolator has each player pin a note to a shared bulletin board pointing at one designated 'captain'; once the captain's own note on the board says 'go,' anyone can walk up to the board later, read the captain's note, and know exactly what to do with their own note — no need to track down the referee.

saying these in an interview costs you the question

  • Thinks Percolator uses a traditional standalone coordinator process like classic 2PC
  • Doesn't know the timestamp oracle's limited, stateless role
  • Confuses the primary row with an ordinary coordinator log
  • Claims Percolator provides full serializable isolation
  • Can't explain how a stale lock gets resolved without a live coordinator

context