Hiring BarSupport

Consistency Models & Partition Behavior

Foundation31 min read5 diagrams

Why This Matters#

Consistency is not a vocabulary quiz. Nobody gets leveled Staff for reciting the difference between sequential and linearizable. Consistency is a question about which anomaly the business can tolerate, what it costs to prevent the others, and who reconciles the damage when the network splits. "Strong consistency" is not a feature you turn on; it's a latency bill paid on every write and an availability bill paid during every partition.

Most candidates reach for one of two defaults. The first: "We'll use a strongly consistent database," with no mention that a cross-region linearizable write costs 60–150ms of consensus per commit and becomes unavailable on the minority side of a partition. The second: "It's eventually consistent," with no bound on how eventual, no statement of what users see in the meantime, and no plan for conflicting writes. Both answers are reasonable. Neither says who pays.

Staff engineers do three things differently. They name the anomaly instead of the model ("a user can post a comment and not see it on refresh" — that's a read-your-writes violation). They scope the guarantee ("linearizable per account, eventual across accounts") instead of choosing one model for the whole system. And they decide partition behavior up front — which operations fail closed, which fail open, and who reconciles afterward — because that's where the CAP theorem stops being a slogan and starts paging people.

The 60-Second Version#

  • Name the anomaly, not the model. Stale read, lost update, write skew, read-your-writes violation, causality violation. Business owners understand anomalies; they don't understand "sequential consistency."
  • Linearizability costs a round trip to a quorum on every write — ~1–5ms within a region, ~60–150ms across regions. It also means the minority side of a partition stops serving writes (and linearizable reads).
  • Most user-facing products need session guarantees, not global linearizability. Read-your-writes + monotonic reads per user covers ~90% of "the user saw stale data" complaints at a fraction of the cost.
  • CAP only applies during a partition. PACELC completes it: else (normal operation) you trade latency vs consistency — and that's the tradeoff you pay every second, not once a year.
  • Quorums: R + W > N gives overlapping reads and writes (e.g., N=3, W=2, R=2), but it is not linearizability on its own — sloppy quorums, hinted handoff, and last-write-wins clocks all leak anomalies.
  • Isolation levels are consistency models for transactions. Postgres defaults to Read Committed; MySQL InnoDB to Repeatable Read. Neither prevents write skew. SERIALIZABLE does, at the cost of aborts you must retry.
  • Last-write-wins silently discards data. With clock skew of even 10–100ms, the "last" write is sometimes the earlier one. If losing a write is unacceptable, you need versions, CRDTs, or a single writer.

How Consistency Models Work#

The Spectrum#

From strongest (most expensive) to weakest (cheapest):

ModelGuarantee (plain English)Typical CostAvailable During Partition?Example Systems
Strict serializabilityTransactions behave as if run one at a time, in real-time orderConsensus + commit wait or timestamp oracleNo (minority side)Spanner, CockroachDB (close), FoundationDB
LinearizabilityEvery single-object op appears to take effect instantly between call and return; everyone sees the same orderQuorum round trip per opNo (minority side)etcd, ZooKeeper writes, single-leader DB reads from leader
Serializability (no real-time)Transactions equivalent to some serial order — not necessarily real-timeLocks or SSI abortsDepends on replicationPostgres SERIALIZABLE on a single node
SequentialAll nodes see the same order; it respects each client's order but not real timeOrdering via a logNoZooKeeper reads (without sync)
CausalIf A could have influenced B, everyone sees A before BDependency tracking (vector clocks, lamport + deps)YesCOPS-style research systems, some geo-replicated stores
Session guaranteesPer-client: read-your-writes, monotonic reads, monotonic writes, writes-follow-readsSticky routing or version tokensYes (mostly)Cosmos DB "Session", MongoDB causal sessions
Bounded stalenessReads lag writes by at most k versions or t secondsReplica lag monitoringYes, until the bound is violatedCosmos DB "Bounded Staleness"
EventualIf writes stop, replicas converge — eventually, no boundCheapestYesDNS, Dynamo-style stores at W=1/R=1, async replicas
Diagram: The Spectrum

Arrows point from stronger to weaker. Everything above causal requires coordination that makes it unavailable on the minority side of a partition; causal and below can stay available. That line — not the CAP acronym — is the practical boundary.

The Four Session Guarantees#

GuaranteeAnomaly It PreventsUser-Visible Symptom Without It
Read-your-writesReading a replica that hasn't seen your own write"I changed my profile picture and it reverted on refresh"
Monotonic readsReading a newer replica, then an older one"The comment appeared, then disappeared, then reappeared"
Monotonic writesYour writes applied out of order"I renamed the doc twice and the first name stuck"
Writes-follow-readsYour write ordered before the write you were responding to"My reply shows up before the message it replies to"

CAP and PACELC#

CAP: during a network Partition, a system must choose between Consistency (linearizability — refuse requests that can't be answered correctly) and Availability (answer every request, possibly with stale or conflicting data).

PACELC: if Partition, choose A or C; Else, choose Latency or Consistency.

System (default config)During Partition (PA/PC)Normal Operation (EL/EC)What That Means
Dynamo-style (Cassandra at ONE, Riak)PAELAlways answers fast; you reconcile conflicts
DynamoDB (eventually consistent reads)PAELStrongly consistent reads available at 2× read cost
Spanner, CockroachDBPCECPays consensus latency on every write; minority side unavailable
Single-leader Postgres + async replicasPC for writesEL for replica readsReplica reads can be stale by the replication lag
MongoDB (majority write concern)PCECWrites wait for majority acknowledgment

Partitions are rare; latency is constant. PACELC's "else" branch is the one you pay for every request, which is why it matters more in design interviews than the CAP branch.

Core Strategies#

Strategy 1: Single Leader, Read from the Leader#

write(key, value):
    leader.write(key, value)                 # leader orders all writes
    replicate_sync_to_quorum(key, value)     # e.g., 1 of 2 followers acks

read_linearizable(key):
    return leader.read(key)                  # plus a lease or read-index check
                                             # so a deposed leader can't answer

When to use: Money, inventory, uniqueness constraints, anything where a stale read causes a wrong decision.

Failure mode: The leader is the throughput ceiling and the failover risk. A leader that doesn't know it's been deposed (GC pause, partition) can serve stale reads unless reads are protected by a lease or a read-index round trip — the classic bug Jepsen has found in more than one database.

Strategy 2: Quorum Reads and Writes#

N = 3 replicas, W = 2, R = 2        # R + W > N → read and write sets overlap

write(key, value, version):
    send to all N; succeed when W acknowledge

read(key):
    ask R replicas; return highest version
    if replicas disagree: read-repair the stale ones

When to use: Leaderless stores where you want tunable consistency per request (Cassandra QUORUM / LOCAL_QUORUM, Riak).

Failure mode: R+W>N is necessary but not sufficient. Sloppy quorums (writes accepted by non-home nodes during failures), concurrent writes resolved by last-write-wins, and a write that reached fewer than W nodes before failing (it may or may not become visible) all produce non-linearizable behavior.

Strategy 3: Session Guarantees via Version Tokens#

write(user, key, value):
    version = primary.write(key, value)      # returns log position / LSN
    return { ok, session_token: version }

read(user, key, session_token):
    replica = pick_replica()
    if replica.applied_version < session_token:
        replica = primary                    # or wait up to 50ms for catch-up
    return replica.read(key)

When to use: Most user-facing products. Users must see their own writes; they rarely notice that other users' writes arrive 500ms late.

Failure mode: Tokens lost across devices or channels. The user writes on the phone and reads on the laptop — a per-device session guarantees nothing across them. Scope the token to the user (stored server-side) if cross-device matters.

Strategy 4: Eventual Consistency with Explicit Conflict Resolution#

ResolutionHowLoses Data?Use For
Last-write-wins (LWW)Highest timestamp winsYes — concurrent writes silently droppedCaches, presence, "last seen"
Version vectors + app mergeKeep siblings; application mergesNo, but app must mergeShopping carts, documents
CRDTsData types whose merge is commutative, associative, idempotentNo (by construction)Counters, sets, collaborative editing
Single writer per keyRoute all writes for a key to one ownerNoPer-user or per-entity state

When to use: Multi-region active-active, offline-capable clients, high-availability writes.

Failure mode: Choosing LWW by default and discovering months later that concurrent edits were silently discarded. LWW is a data-loss policy — someone in product must sign off on it.

Strategy 5: Pick the Transaction Isolation Level on Purpose#

Isolation LevelPreventsStill AllowsDefault In
Read Uncommitted(almost nothing)Dirty readsRarely used
Read CommittedDirty reads, dirty writesNon-repeatable reads, read skew, lost updates, write skew, phantomsPostgreSQL, Oracle, SQL Server
Repeatable Read / Snapshot Isolation+ non-repeatable reads, read skew (lost updates in Postgres RR)Write skew (and phantoms in some engines)MySQL InnoDB (RR)
SerializableAll of the aboveNothing — but transactions abort and must retryCockroachDB, FoundationDB; opt-in in Postgres (SSI)

When to use Serializable: invariants that span multiple rows — "at least one doctor on call," "total allocations ≤ budget," "username unique across two tables."

Failure mode: Turning on Serializable without retry logic. Under contention, serialization failures (40001 in Postgres) can reach several percent of transactions; code that doesn't retry turns them into user-visible errors.

Partition Behavior: The Hard Sub-Problem#

CAP says you must choose during a partition. It doesn't say what to choose, per operation, or who cleans up. That's the design work.

Decide Per Operation, Not Per System#

OperationFail Closed (C) or Open (A)?WhyWho Reconciles After
Charge a card / move moneyClosedA double-spend costs real money and regulator attentionNobody — prevented
Reserve the last seatClosed (or bounded oversell)Overselling means a human apologizes at the gateOps team, if oversell is allowed
Add item to cartOpenA lost "add" loses a sale; a duplicate is harmlessMerge on reconnect (union of items)
Post a commentOpenDelay is fine; failure feels brokenAsync replication catches up
Change password / revoke sessionClosed for the change, open for reads with short TTLSecurity change must be authoritativeSecurity team, if revocation lags
Like counterOpenOff-by-a-few is invisibleCRDT counter converges
Username / email uniquenessClosedTwo accounts with one email is a support ticket foreverPrevented

🎯 Staff Move: "I'm not going to pick CP or AP for the whole system. Payments and seat holds fail closed — if the region can't reach a quorum, those endpoints return 503 and the UI says 'try again.' Cart, comments, and likes fail open and merge on reconnect. Product has to sign off on the list, because they're choosing which failures customers see."

What a Partition Actually Looks Like#

Region A (2 of 3 replicas)            |  Region B (1 of 3 replicas)
--------------------------------------+--------------------------------------
Majority side: elects/keeps leader    |  Minority side: loses quorum
CP ops: continue normally             |  CP ops: time out → 503 after ~1–5s
AP ops: continue                      |  AP ops: accept writes locally
                                      |  Clients: see their own writes; not A's
--- partition heals after 4 minutes ---
Replicas exchange missed writes. CP data: B catches up from A's log.
AP data: concurrent writes to the same key now conflict → merge or LWW.

The partition's real cost is paid after it heals: conflicting writes, duplicated side effects (emails sent twice), and derived data (search indexes, counters, caches) that saw one side's history.

The Anomalies, Concretely#

AnomalyExamplePrevented By
Stale readUser sees yesterday's balance from a lagging replicaLeader reads, bounded staleness, version tokens
Lost updateTwo clients read stock=5, both write stock=4; one sale vanishesSELECT … FOR UPDATE, compare-and-set, atomic decrement
Read skewTransfer reads account A before and account B after a concurrent transfer; total looks wrongSnapshot isolation
Write skewTwo doctors each check "another doctor is on call," both go off callSerializable, or materialize the conflict (lock a shared row)
PhantomA range check (COUNT(*) WHERE room=7 AND overlaps) misses a concurrently inserted bookingSerializable, predicate/range locks, exclusion constraints
Causality violationA reply appears before the message it replies toCausal consistency, writes-follow-reads
Split brainTwo leaders accept conflicting writesConsensus with fencing tokens / epochs
Full reasoning: write skew, the anomaly snapshot isolation misses

Snapshot isolation gives each transaction a consistent snapshot and aborts one of two transactions that write the same row. Write skew happens when two transactions read an overlapping set, then write different rows based on what they read:

Invariant: at least 1 doctor on call for shift 42.
T1: SELECT count(*) FROM oncall WHERE shift=42   → 2
T2: SELECT count(*) FROM oncall WHERE shift=42   → 2
T1: DELETE FROM oncall WHERE shift=42 AND doctor='alice'
T2: DELETE FROM oncall WHERE shift=42 AND doctor='bob'
Both commit under Snapshot Isolation → 0 doctors on call.

No row was written by both, so SI's write-write conflict check never fires. Fixes, in order of preference:

  1. Serializable isolation (Postgres SSI detects the read-write dependency cycle and aborts one transaction).
  2. Materialize the conflict: SELECT … FROM shifts WHERE id=42 FOR UPDATE so both transactions contend on one row.
  3. A constraint the database enforces (exclusion constraint, unique index) so the invariant can't be violated regardless of isolation.

The interview signal is recognizing write skew when the interviewer describes a booking or on-call invariant — and not claiming "we use transactions" fixes it.

Failure Scenario: The Async Replica Promotion#

t=0      Primary in us-east-1 serving 8K writes/s; async replica in us-west-2, lag ~200ms.
t=+0s    us-east-1 primary host fails hard (disk controller).
t=+30s   Automated failover promotes the us-west-2 replica.
t=+31s   ~1,600 committed writes (200ms × 8K/s) exist only on the dead primary.
t=+5min  Customers report orders that "disappeared"; some were charged via the payment provider.
t=+2h    Old primary's disk recovered; 1,600 orphaned writes identified.
t=+2d    Manual reconciliation: re-insert non-conflicting rows; refund 37 double-charged orders.

Detection: replication.lag.seconds before failover (it was known); post-failover diff of the old primary's WAL vs the new primary; spike in orders.not_found support tickets. Blast radius: every write committed in the last replication-lag window — here ~1,600 orders. Mitigation: fence the old primary so it can't accept writes if it comes back; reconcile from its WAL. Prevention: semi-synchronous replication to at least one replica before ack (≈ +1–2ms in-region), or make cross-region failover a human decision with an explicit "accept data loss of N seconds" sign-off. Owner: database platform team owns failover policy; the payments team owns the reconciliation runbook, because only they know which lost writes had external side effects.

Operational Reality Matrix#

FailureDetection SignalBlast RadiusMitigationOwner
Replica lag spikereplication.lag.seconds > SLO (e.g., 1s)Every stale-read pathRoute session reads to primary; shed replica trafficDB platform
Failover with data lossWrites missing post-promotion; WAL diffLast lag-window of writesFence old primary; reconcileDB platform + domain team
Split brainTwo nodes report leader; leader.epoch divergenceAll writes during overlapFencing tokens; kill old leaderConsensus/infra owner
LWW discarding writesconflict.lww_overwrites counter; user reportsConcurrent editorsMove to versions/CRDT for that entityProduct team owning the entity
Serialization-failure stormdb.tx.serialization_failures > 5%Contended invariantsRetries with backoff; narrow transactionsService owner
Minority-side unavailabilityCP endpoints 503 in one regionUsers routed to that regionRoute users to majority regionSRE / traffic team

Visual Guide#

Choosing a Consistency Level per Operation#

Diagram: Choosing a Consistency Level per Operation

Read-Your-Writes with a Version Token#

Diagram: Read-Your-Writes with a Version Token

Behavior During a Partition#

Diagram: Behavior During a Partition

Implementation Patterns#

Fencing Tokens#

Every leader term gets a monotonically increasing epoch. Storage rejects writes carrying an epoch lower than the highest it has seen. A deposed leader waking from a 30-second GC pause can't corrupt data because its token is stale.

Consistency Scoped by Key#

Most systems need strong consistency per entity, not globally. Route all writes for an account to one partition leader; you get linearizability for that account and horizontal scale across accounts. Cross-entity invariants ("transfer between two accounts") then need a transaction protocol or a saga — and that's where the real cost lives.

Commit Wait and Hybrid Clocks#

Spanner waits out TrueTime's uncertainty (ε, typically a few milliseconds) before acknowledging a commit, so timestamp order matches real-time order. CockroachDB uses hybrid logical clocks with a bounded max clock offset (500ms default) and restarts reads that encounter uncertainty. Both convert clock uncertainty into latency, which is the right trade: you pay a few ms instead of getting anomalies.

Outbox for Cross-System Consistency#

A database commit plus a Kafka publish is two systems and no transaction. Write the event into an outbox table in the same transaction; a relay publishes it. Consumers see events at-least-once and in commit order per key. This gives you eventual consistency with no lost events across service boundaries — which is what most "consistency between services" questions are actually asking for.

Quorum Configurations#

NWRTolerates for WritesTolerates for ReadsUse
3221 down1 downDefault for critical data
3310 down2 downRead-heavy, rare writes (writes block on any failure)
3132 down0 downWrite-heavy, reads can wait (rare)
3112 down2 downNo overlap — eventual; caches, metrics
5332 down2 downSurvive an AZ plus one node
6 (3 per region)LOCAL_QUORUM (2)LOCAL_QUORUM (2)1 per region1 per regionMulti-region, strong within a region, async across

When NOT to Pay for Strong Consistency#

  • Derived or recomputable data — search indexes, recommendation scores, analytics rollups. Rebuild beats coordinate.
  • Data whose staleness no human can perceive — like counts, view counts, "last seen 2 minutes ago."
  • Reads that inform, not decide — the product page stock count; the reservation decides.
  • Cross-region paths on the user's critical latency budget — 70–150ms per write is usually a product regression the business didn't ask for.

The Numbers in Context#

NumberValueWhat It Means for Your Design
In-region quorum write~1–5msStrong consistency inside a region is cheap. Default to it for critical data.
Cross-region quorum write~60–150ms (one round trip to the nearest majority)Every linearizable global write pays this. Budget it or avoid it.
Async replica lag (healthy)~10ms–1sSession guarantees matter; users refresh within 1s.
Async replica lag (unhealthy)seconds to hoursReplica reads need a staleness bound and an alert.
Data loss on async failoverlag × write rate (200ms × 8K/s = 1,600 writes)Async cross-region failover is a data-loss decision.
Spanner TrueTime εtypically low single-digit msCommit wait is small but non-zero on every write.
CockroachDB max clock offset500ms defaultNodes whose clocks drift past this are shut down.
DynamoDB strongly consistent read2× the read capacity of an eventual readStrong reads double the read bill.
Clock skew with NTP~1–10ms typical, 100ms+ during incidentsLWW with wall clocks will order some writes wrong.
Serializable abort rate<1% typical; several % under contentionRetry logic is part of the design, not an afterthought.
Quorum configN=3, W=2, R=2Survives 1 replica down for both reads and writes.
Raft/Paxos group3 or 5 voters5 survives 2 failures at the cost of a slower majority.

How This Shows Up in Interviews#

Scenario 1: "Is your system CP or AP?"#

Don't answer with one letter. "Per operation. Balance changes and seat holds are CP — they need a quorum and fail with a retryable error on the minority side. Feed, likes, and comments are AP with session guarantees so users see their own writes. And in normal operation — the PACELC 'else' — I'm choosing latency for reads by serving from local replicas with a 1-second staleness bound, and consistency for writes by going to the leader."

Scenario 2: "Design the inventory service for a flash sale across two regions." (Full Walkthrough)#

Step 1 — Name the anomaly that matters. "The anomaly the business can't tolerate is overselling: two buyers both getting the last unit. Stale reads of the stock count on the product page are fine — nobody is harmed if the page says '12 left' and it's really 9."

Step 2 — Scope the strong guarantee. "So I split reads from reservations. Product-page stock counts are eventually consistent, served from a cache with a 2-second TTL. The reservation — decrementing stock — is linearizable per SKU."

Step 3 — Place the leader. "Cross-region consensus on every reservation costs ~70ms and makes the minority region unable to sell during a partition. Instead each SKU's inventory has a home region leader, and I pre-allocate stock: region A gets 60% of units, region B 40%, based on forecast demand. Each region decrements its own allocation locally — linearizable within the region at ~2ms."

Step 4 — Rebalance allocations. "When one region runs low, it requests more units from the other through a single transfer transaction. That's the only cross-region consistent operation, and it's rare."

Step 5 — Partition behavior. "If the regions partition, each keeps selling from its own allocation — no oversell is possible, because neither region can spend units it doesn't hold. The worst case is that region A sells out while B still has stock. Product signs off that 'sold out a little early in one region' is acceptable during a partition; overselling is not."

Step 6 — Owners. "Inventory team owns allocation policy and the transfer protocol; the flash-sale team owns the demand forecast that sets initial splits."

Why this is a Staff answer: It identifies the one anomaly that matters, applies strong consistency only where that anomaly lives, converts a global consensus problem into local ones through allocation, and gets product to sign off on partition behavior in business terms.

Scenario 3: "Users say their comment disappears after posting, then comes back."#

This is a monotonic-reads plus read-your-writes violation. "The write went to the primary; the next read hit a lagging replica, the one after that hit a caught-up replica. I'd return an LSN session token on write and route reads to replicas that have applied it — or pin the user to the primary for ~5 seconds after a write. I'd also add replication.lag.seconds per replica to the load balancer's health check so a replica lagging past 1s stops receiving reads."

Scenario 4: "Two admins edit the same config simultaneously. One change is lost."#

Lost update. "Every config write carries the version it read — UPDATE config SET … , version = v+1 WHERE id = ? AND version = v. If zero rows update, the client re-reads and shows the conflict. For a config system, I'd rather make the second admin merge explicitly than silently pick a winner."

Advanced Patterns#

PatternHow It WorksWhen to Use
Leader leases / read indexLeader proves it's still leader before serving a linearizable readLinearizable reads without a full write round trip
Follower reads with a timestampRead from a local follower "as of" a slightly past timestamp (e.g., −5s)Low-latency consistent snapshots in multi-region DBs (Spanner/CockroachDB)
Escrow / allocationSplit a numeric budget across regions; each spends locallyInventory, quotas, rate limits across regions
CRDTsG-Counter, PN-Counter, OR-Set, LWW-Register, sequence CRDTsCounters, sets, collaborative text, offline-first apps
Sagas with compensationCross-service invariants maintained by forward steps + undo stepsOrder → payment → shipping across services
Transactional outbox + CDCEvents committed atomically with state; published from the logKeeping search, cache, and analytics consistent with the DB
Consistency-level per requestClient specifies ONE / LOCAL_QUORUM / strong per callMixed workloads on one tunable store

The Principal Lens#

Why L7 Sees This Problem Differently#

A Staff engineer picks the right consistency model for each operation in their system. A Principal engineer notices that the anomalies users actually hit rarely live inside one database — they live between services. The order service is linearizable, the payment service is linearizable, and the customer still sees "order confirmed, payment failed" because the two are stitched together by an async event with no contract about ordering, duplication, or staleness. At L7, consistency becomes a cross-team contract: what each service promises about the freshness and ordering of the data it exposes, and which team owns reconciliation when those promises are broken.

The Org-Level Fault Line#

One consistency posture for the whole company vs per-service choice.

OptionWhat WorksWhat BreaksWho Pays
Mandate a strongly consistent distributed DB everywhereFewer anomaly classes; simpler reasoning for product teamsEvery write pays consensus latency; database bill 2–5× higher; single vendor dependencyLatency-sensitive teams; the infrastructure budget
Every team chooses freelyEach service tuned to its workloadInconsistent guarantees at every boundary; reconciliation bugs nobody ownsCustomer support, finance reconciliation, on-call at 3am
Tiered posture with published contractsStrong where invariants live; eventual elsewhere; boundaries documentedRequires a classification process and enforcement in design reviewPlatform/architecture team maintaining the tiers

The Principal position: a small number of consistency tiers (e.g., Ledger, Transactional, Session, Eventual), each with a paved-road datastore and a published staleness/ordering contract, and a rule that every cross-service data flow declares which tier's contract it offers.

Cost Model#

Assumptions: a regional strongly consistent SQL cluster (3 replicas) ≈ 3× the cost of a single primary with async replicas; a multi-region consensus database ≈ 1.5–3× the regional cluster (more replicas, cross-region transfer, premium pricing); cross-region transfer ~$0.02/GB; engineer ~$25K/month.

ScalePostureInfra $/monthReconciliation / Incident CostHeadcount
Startup (1 region, 1 Postgres primary + 2 replicas)Strong in-region, session reads~$3–8K~0 — few anomalies, one DB0 dedicated
Growth (1 region, 30 services, 10 DBs, async events)Per-service choice, outbox for events~$60–120K1 FTE-equivalent/year on reconciliation scripts and incident cleanup ($25K/month)1–2 on data platform
Large (3 regions, 300 services)Tiered: Ledger tier on multi-region consensus DB, the rest regional + async~$1–2M; the Ledger tier is ~10% of data and ~30–40% of DB spendWithout tiers: finance reconciliation team of 3–5; with tiers: ~1–26–10 on data platform

The Principal insight: making everything globally consistent roughly doubles the database bill; making nothing consistent staffs a reconciliation team. The tiered model spends consensus dollars only where an anomaly costs more than the consensus.

The 3-Year Evolution Path#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionDoorReversal Cost
Isolation level on a single DBTwo-wayConfig + retry logic; days
Session guarantees via tokensTwo-wayClient/API change; weeks
Last-write-wins for a user-data entityOne-way (for the data lost)Discarded writes are unrecoverable
Active-active multi-writer for an entityMostly one-wayConflict semantics leak into every client and downstream consumer
Public API promising read-your-writes or orderingOne-wayExternal clients depend on it forever
Choice of a globally consistent DB vendorMostly one-waySQL dialect, transaction semantics, and pricing lock-in; multi-year migration

The Standard I'd Write#

RFC: Data Consistency Tiers (v1)

Scope: All persistent data owned by production services, and every cross-service data flow (APIs, events, CDC).

Tiers: Ledger (strict serializable, multi-region, zero data loss on failover) · Transactional (serializable or linearizable within a region; RPO ≤ 1s cross-region) · Session (read-your-writes + monotonic reads per user; replica staleness ≤ 1s p99) · Eventual (converges within a documented bound; conflicts resolved by a declared policy).

MUST:

  1. Every datastore and every event stream declares its tier in the service catalog.
  2. Any data involving money movement or legal records MUST be Ledger tier.
  3. Last-write-wins MUST NOT be used for user-authored content without a product sign-off recorded in the design doc.
  4. Cross-service events MUST be published via transactional outbox or CDC; dual writes are prohibited.
  5. Services MUST emit replication.lag.seconds and consistency.violation metrics (e.g., RYW token misses).

SHOULD: prefer scoping strong guarantees per entity over global ones; prefer escrow/allocation over cross-region consensus for numeric budgets.

Exceptions: Architecture review; exceptions carry a named reconciliation owner.

Success metrics: reconciliation incidents per quarter trending to zero; 100% of Ledger data on the Ledger tier within 4 quarters; no dual-write findings in design review.

What I'd Tell the VP#

"Some of our data has to be exactly right — money, inventory, legal records — and some of it can be a second or two behind without anyone noticing. Today every team decides for itself, so we pay for perfect accuracy in places that don't need it and we get mismatches in places that do; finance spent six weeks last year reconciling one of them. We want four clear categories, a standard database for each, and a rule that money always goes in the strictest one. That puts our most expensive database spend on about a tenth of our data instead of all of it. It also shrinks the reconciliation work that currently takes several people."

Principal Interview Signals#

SignalWhat It Sounds Like
Moves the problem to service boundaries"Both databases are linearizable; the anomaly lives in the event between them."
Prices consistency"Global consensus for all 300 services doubles the DB bill. For the 10% of data that's ledger, it's cheap insurance."
Treats LWW as a policy decision"Last-write-wins is a data-loss policy. I want product's signature on it, per entity."
Defines tiers and contracts"Every event stream declares its ordering and staleness contract, and the consumer designs to that contract."
Knows the one-way doors"Once we publish read-your-writes in our public API, we can never go back to plain async replicas for that endpoint."

Staff answers that L7 interviewers find insufficient:

  • "We'll use Spanner so we don't have to think about consistency." — Buys correctness inside one database and ignores the boundaries where anomalies actually occur.
  • "Payments is CP and feed is AP." — Right classification, no mechanism to make 300 services classify consistently.
  • "We'll reconcile with a nightly job." — Names a mechanism with no owner, no SLA, and no plan for side effects that already escaped.

In the Wild#

These are public, documented examples.

Google Spanner: External Consistency via TrueTime#

Spanner (OSDI 2012) provides externally consistent (strictly serializable) transactions across globally distributed data. It uses TrueTime — GPS receivers and atomic clocks in each data center — to bound clock uncertainty, and makes each commit wait out that uncertainty before acknowledging. Replication within each shard is Paxos-based, so writes pay a majority round trip in addition to commit wait.

Staff insight: Spanner doesn't make consistency free; it makes the cost explicit and bounded — a few milliseconds of commit wait plus consensus latency. When you cite it, cite the price: "strict serializability at the cost of a majority round trip on every write."

Amazon Dynamo: Availability First, Merge Later#

The 2007 Dynamo paper described an always-writable store for the shopping cart. It used sloppy quorums and hinted handoff to accept writes during failures, and vector clocks to detect conflicting versions, which were returned to the application to merge. For the cart, the merge was a union of items — which meant deleted items could reappear after a conflict, a known and accepted anomaly.

Staff insight: Amazon chose availability for "add to cart" because a rejected add loses a sale, and accepted a resurrected item as the cost. That's exactly the "name the anomaly, name who pays" reasoning: the customer occasionally removes an item twice; the business never loses an add.

Jepsen: Advertised vs Actual Guarantees#

Kyle Kingsbury's Jepsen project has published analyses of many distributed databases and queues under partitions, clock skew, and process pauses. Repeatedly, it has found systems that violated their documented consistency or isolation guarantees in specific configurations — lost acknowledged writes, stale reads from deposed leaders, and anomalies prohibited by the claimed isolation level. Vendors have fixed many of these issues in response.

Staff insight: A vendor's consistency claim is a hypothesis, not a guarantee. The Staff move is to name the configuration you rely on (write concern, read concern, isolation level) and how you'd test it — not to quote the marketing page.


Staff Calibration#

What Staff Engineers Say (That Seniors Don't)#

ConceptSenior (L5)Staff (L6)Principal (L7)
Choosing a model"We need strong consistency""Linearizable per account for balance changes; session guarantees for profile reads; eventual for counters""Four company-wide tiers with published contracts; every datastore and event stream declares one"
CAP"It's a CP system""CP for payments and seat holds, AP for cart and comments — and in normal operation I'm trading latency for consistency per PACELC""Partition behavior is a product policy table that product signs off on, reviewed yearly"
Stale reads"Read from the primary""Version tokens route only the writing user to fresh replicas; everyone else reads local replicas with a 1s bound""RYW token misses are an SLO with an owner, not a bug report"
Transactions"Wrap it in a transaction""Read Committed won't stop write skew on this invariant; I'll use Serializable with retries or lock the shift row""Invariants that span services can't be transactions; they're sagas with a named reconciliation owner"
Conflicts"Last write wins""LWW silently drops concurrent edits; carts merge, counters are CRDTs, documents keep versions""LWW on user data requires recorded product sign-off; it's a data-loss policy"
Failover"Promote the replica""Async promotion loses lag × write-rate writes — 1,600 here; fence the old primary and reconcile""Cross-region failover RPO is a business decision set per tier, rehearsed in game days"
Why "Transactions" separates levels

"Use a transaction" is not wrong — it's underspecified. Senior candidates assume a transaction means "correct." Staff candidates know the default isolation level (Read Committed in Postgres) and can name the anomaly it permits for the specific invariant. Principal candidates recognize that the hardest invariants in a microservice architecture span services, where no single database transaction exists, and design the reconciliation and ownership model around that fact.

Why "Failover" separates levels

Promoting a replica is a routine operation that silently decides how much data you lose. Staff engineers compute it (lag × write rate) and fence the old leader. Principal engineers turn the number into a per-tier RPO that the business approves and that is exercised in game days, so the decision isn't made for the first time at 3am.

Staff Sentence Templates#

"The anomaly the business can't tolerate here is [anomaly], so [operation] is [guarantee]; everything else is [weaker guarantee] with a [bound] staleness bound."

"During a partition, [operation] fails closed and returns [error]; [other operation] accepts writes locally and reconciles by [merge policy]. Product signed off on [user-visible effect]."

"Async failover loses [lag] × [write rate] ≈ [N] writes. [Team] owns reconciliation, because only they know which of those had external side effects."

Key Takeaway: Consistency is priced per operation. Buy linearizability where an anomaly costs more than a quorum round trip; buy session guarantees where users would notice; let everything else be eventual with a stated bound.

Common Interview Traps#

  • Choosing one model for the whole system. Scope guarantees per operation or per entity.
  • Treating R + W > N as linearizable. Sloppy quorums, LWW, and partial writes leak anomalies.
  • Saying "eventually consistent" without a bound. How eventual? 100ms? An hour? What do users see meanwhile?
  • Assuming the default isolation level prevents write skew. Read Committed and Snapshot Isolation don't.
  • Using wall-clock LWW across regions. Clock skew reorders writes; you lose data silently.
  • Dual writes to DB and queue. One succeeds, one fails, and they diverge. Use an outbox.
  • Ignoring the healed partition. The hard part is reconciliation and side effects, not the partition itself.
  • Trusting vendor claims without naming the configuration. Write concern and read concern matter.

Practice Drill#

Prompt: "We're expanding from one region to two, active-active. Product wants users to be able to edit their shared documents from either region with low latency. Walk me through consistency."

Staff Answer

I'd start by splitting the data by anomaly tolerance, because "active-active" means different things for different entities. Document content is the hard one: two users in two regions can edit the same document concurrently, so either every edit pays ~70–90ms of cross-region consensus (latency product rejected) or edits are accepted locally and merged. Merging text safely needs a CRDT or OT-based sequence type, not LWW — LWW would silently discard one user's paragraph. So document content becomes a CRDT, replicated asynchronously between regions, with each region applying remote operations as they arrive (typically sub-second). Document permissions and sharing are different: revoking access must be authoritative, so ACL changes go through a single home region per document (linearizable), and reads use a short cache (~5s) with explicit invalidation — product signs off that a revoked user might see up to ~5s more of edits. Billing/quota (number of documents per plan) uses escrow: each region gets a share of the quota and rebalances occasionally. During a partition, content editing continues in both regions and merges on heal; ACL changes for documents homed in the unreachable region fail with a clear error; quota continues from local allocation. Metrics: crdt.merge.lag.seconds between regions, acl.cache.staleness, and consistency.violation for any RYW miss. Owners: the collaboration team owns the CRDT and merge semantics; the identity/permissions team owns ACL authority and the revocation SLA.

Why this is L6:

  • Splits one "active-active" requirement into three entities with different anomaly tolerances.
  • Rejects LWW for user content and names the data loss it causes.
  • Defines partition behavior per entity and gets product sign-off on the revocation window.
  • Attaches metrics and owners to each guarantee.

What L7 adds:

  • Recognizes that CRDT-based active-active is a multi-year platform commitment (storage format, client libraries, offline support) and decides whether to build it once for the company rather than for this product.
  • Prices the second region — duplicated infra plus cross-region replication traffic — against the latency gain for the affected user base, and presents it as a business decision.
  • Adds the document entity to the org's consistency-tier catalog so downstream consumers (search, analytics, exports) design to its eventual contract.

Where This Appears#

Related Foundations & Patterns: Consistent Hashing · Sharding & Partitioning · Dealing with Contention · Distributed Coordination

Related Technologies: PostgreSQL · Cassandra · DynamoDB · ZooKeeper & etcd

  1. Loading the index…