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.
SERIALIZABLEdoes, 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):
| Model | Guarantee (plain English) | Typical Cost | Available During Partition? | Example Systems |
|---|---|---|---|---|
| Strict serializability | Transactions behave as if run one at a time, in real-time order | Consensus + commit wait or timestamp oracle | No (minority side) | Spanner, CockroachDB (close), FoundationDB |
| Linearizability | Every single-object op appears to take effect instantly between call and return; everyone sees the same order | Quorum round trip per op | No (minority side) | etcd, ZooKeeper writes, single-leader DB reads from leader |
| Serializability (no real-time) | Transactions equivalent to some serial order — not necessarily real-time | Locks or SSI aborts | Depends on replication | Postgres SERIALIZABLE on a single node |
| Sequential | All nodes see the same order; it respects each client's order but not real time | Ordering via a log | No | ZooKeeper reads (without sync) |
| Causal | If A could have influenced B, everyone sees A before B | Dependency tracking (vector clocks, lamport + deps) | Yes | COPS-style research systems, some geo-replicated stores |
| Session guarantees | Per-client: read-your-writes, monotonic reads, monotonic writes, writes-follow-reads | Sticky routing or version tokens | Yes (mostly) | Cosmos DB "Session", MongoDB causal sessions |
| Bounded staleness | Reads lag writes by at most k versions or t seconds | Replica lag monitoring | Yes, until the bound is violated | Cosmos DB "Bounded Staleness" |
| Eventual | If writes stop, replicas converge — eventually, no bound | Cheapest | Yes | DNS, Dynamo-style stores at W=1/R=1, async replicas |
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#
| Guarantee | Anomaly It Prevents | User-Visible Symptom Without It |
|---|---|---|
| Read-your-writes | Reading a replica that hasn't seen your own write | "I changed my profile picture and it reverted on refresh" |
| Monotonic reads | Reading a newer replica, then an older one | "The comment appeared, then disappeared, then reappeared" |
| Monotonic writes | Your writes applied out of order | "I renamed the doc twice and the first name stuck" |
| Writes-follow-reads | Your 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) | PA | EL | Always answers fast; you reconcile conflicts |
| DynamoDB (eventually consistent reads) | PA | EL | Strongly consistent reads available at 2× read cost |
| Spanner, CockroachDB | PC | EC | Pays consensus latency on every write; minority side unavailable |
| Single-leader Postgres + async replicas | PC for writes | EL for replica reads | Replica reads can be stale by the replication lag |
| MongoDB (majority write concern) | PC | EC | Writes 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#
| Resolution | How | Loses Data? | Use For |
|---|---|---|---|
| Last-write-wins (LWW) | Highest timestamp wins | Yes — concurrent writes silently dropped | Caches, presence, "last seen" |
| Version vectors + app merge | Keep siblings; application merges | No, but app must merge | Shopping carts, documents |
| CRDTs | Data types whose merge is commutative, associative, idempotent | No (by construction) | Counters, sets, collaborative editing |
| Single writer per key | Route all writes for a key to one owner | No | Per-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 Level | Prevents | Still Allows | Default In |
|---|---|---|---|
| Read Uncommitted | (almost nothing) | Dirty reads | Rarely used |
| Read Committed | Dirty reads, dirty writes | Non-repeatable reads, read skew, lost updates, write skew, phantoms | PostgreSQL, 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) |
| Serializable | All of the above | Nothing — but transactions abort and must retry | CockroachDB, 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#
| Operation | Fail Closed (C) or Open (A)? | Why | Who Reconciles After |
|---|---|---|---|
| Charge a card / move money | Closed | A double-spend costs real money and regulator attention | Nobody — prevented |
| Reserve the last seat | Closed (or bounded oversell) | Overselling means a human apologizes at the gate | Ops team, if oversell is allowed |
| Add item to cart | Open | A lost "add" loses a sale; a duplicate is harmless | Merge on reconnect (union of items) |
| Post a comment | Open | Delay is fine; failure feels broken | Async replication catches up |
| Change password / revoke session | Closed for the change, open for reads with short TTL | Security change must be authoritative | Security team, if revocation lags |
| Like counter | Open | Off-by-a-few is invisible | CRDT counter converges |
| Username / email uniqueness | Closed | Two accounts with one email is a support ticket forever | Prevented |
🎯 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#
| Anomaly | Example | Prevented By |
|---|---|---|
| Stale read | User sees yesterday's balance from a lagging replica | Leader reads, bounded staleness, version tokens |
| Lost update | Two clients read stock=5, both write stock=4; one sale vanishes | SELECT … FOR UPDATE, compare-and-set, atomic decrement |
| Read skew | Transfer reads account A before and account B after a concurrent transfer; total looks wrong | Snapshot isolation |
| Write skew | Two doctors each check "another doctor is on call," both go off call | Serializable, or materialize the conflict (lock a shared row) |
| Phantom | A range check (COUNT(*) WHERE room=7 AND overlaps) misses a concurrently inserted booking | Serializable, predicate/range locks, exclusion constraints |
| Causality violation | A reply appears before the message it replies to | Causal consistency, writes-follow-reads |
| Split brain | Two leaders accept conflicting writes | Consensus 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:
- Serializable isolation (Postgres SSI detects the read-write dependency cycle and aborts one transaction).
- Materialize the conflict:
SELECT … FROM shifts WHERE id=42 FOR UPDATEso both transactions contend on one row. - 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#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Replica lag spike | replication.lag.seconds > SLO (e.g., 1s) | Every stale-read path | Route session reads to primary; shed replica traffic | DB platform |
| Failover with data loss | Writes missing post-promotion; WAL diff | Last lag-window of writes | Fence old primary; reconcile | DB platform + domain team |
| Split brain | Two nodes report leader; leader.epoch divergence | All writes during overlap | Fencing tokens; kill old leader | Consensus/infra owner |
| LWW discarding writes | conflict.lww_overwrites counter; user reports | Concurrent editors | Move to versions/CRDT for that entity | Product team owning the entity |
| Serialization-failure storm | db.tx.serialization_failures > 5% | Contended invariants | Retries with backoff; narrow transactions | Service owner |
| Minority-side unavailability | CP endpoints 503 in one region | Users routed to that region | Route users to majority region | SRE / traffic team |
Visual Guide#
Choosing a Consistency Level per Operation#
Read-Your-Writes with a Version Token#
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#
| N | W | R | Tolerates for Writes | Tolerates for Reads | Use |
|---|---|---|---|---|---|
| 3 | 2 | 2 | 1 down | 1 down | Default for critical data |
| 3 | 3 | 1 | 0 down | 2 down | Read-heavy, rare writes (writes block on any failure) |
| 3 | 1 | 3 | 2 down | 0 down | Write-heavy, reads can wait (rare) |
| 3 | 1 | 1 | 2 down | 2 down | No overlap — eventual; caches, metrics |
| 5 | 3 | 3 | 2 down | 2 down | Survive an AZ plus one node |
| 6 (3 per region) | LOCAL_QUORUM (2) | LOCAL_QUORUM (2) | 1 per region | 1 per region | Multi-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#
| Number | Value | What It Means for Your Design |
|---|---|---|
| In-region quorum write | ~1–5ms | Strong 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–1s | Session guarantees matter; users refresh within 1s. |
| Async replica lag (unhealthy) | seconds to hours | Replica reads need a staleness bound and an alert. |
| Data loss on async failover | lag × write rate (200ms × 8K/s = 1,600 writes) | Async cross-region failover is a data-loss decision. |
| Spanner TrueTime ε | typically low single-digit ms | Commit wait is small but non-zero on every write. |
| CockroachDB max clock offset | 500ms default | Nodes whose clocks drift past this are shut down. |
| DynamoDB strongly consistent read | 2× the read capacity of an eventual read | Strong reads double the read bill. |
| Clock skew with NTP | ~1–10ms typical, 100ms+ during incidents | LWW with wall clocks will order some writes wrong. |
| Serializable abort rate | <1% typical; several % under contention | Retry logic is part of the design, not an afterthought. |
| Quorum config | N=3, W=2, R=2 | Survives 1 replica down for both reads and writes. |
| Raft/Paxos group | 3 or 5 voters | 5 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#
| Pattern | How It Works | When to Use |
|---|---|---|
| Leader leases / read index | Leader proves it's still leader before serving a linearizable read | Linearizable reads without a full write round trip |
| Follower reads with a timestamp | Read from a local follower "as of" a slightly past timestamp (e.g., −5s) | Low-latency consistent snapshots in multi-region DBs (Spanner/CockroachDB) |
| Escrow / allocation | Split a numeric budget across regions; each spends locally | Inventory, quotas, rate limits across regions |
| CRDTs | G-Counter, PN-Counter, OR-Set, LWW-Register, sequence CRDTs | Counters, sets, collaborative text, offline-first apps |
| Sagas with compensation | Cross-service invariants maintained by forward steps + undo steps | Order → payment → shipping across services |
| Transactional outbox + CDC | Events committed atomically with state; published from the log | Keeping search, cache, and analytics consistent with the DB |
| Consistency-level per request | Client specifies ONE / LOCAL_QUORUM / strong per call | Mixed 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.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Mandate a strongly consistent distributed DB everywhere | Fewer anomaly classes; simpler reasoning for product teams | Every write pays consensus latency; database bill 2–5× higher; single vendor dependency | Latency-sensitive teams; the infrastructure budget |
| Every team chooses freely | Each service tuned to its workload | Inconsistent guarantees at every boundary; reconciliation bugs nobody owns | Customer support, finance reconciliation, on-call at 3am |
| Tiered posture with published contracts | Strong where invariants live; eventual elsewhere; boundaries documented | Requires a classification process and enforcement in design review | Platform/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.
| Scale | Posture | Infra $/month | Reconciliation / Incident Cost | Headcount |
|---|---|---|---|---|
| Startup (1 region, 1 Postgres primary + 2 replicas) | Strong in-region, session reads | ~$3–8K | ~0 — few anomalies, one DB | 0 dedicated |
| Growth (1 region, 30 services, 10 DBs, async events) | Per-service choice, outbox for events | ~$60–120K | 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 spend | Without tiers: finance reconciliation team of 3–5; with tiers: ~1–2 | 6–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#
One-Way Doors vs Two-Way Doors#
| Decision | Door | Reversal Cost |
|---|---|---|
| Isolation level on a single DB | Two-way | Config + retry logic; days |
| Session guarantees via tokens | Two-way | Client/API change; weeks |
| Last-write-wins for a user-data entity | One-way (for the data lost) | Discarded writes are unrecoverable |
| Active-active multi-writer for an entity | Mostly one-way | Conflict semantics leak into every client and downstream consumer |
| Public API promising read-your-writes or ordering | One-way | External clients depend on it forever |
| Choice of a globally consistent DB vendor | Mostly one-way | SQL 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:
- Every datastore and every event stream declares its tier in the service catalog.
- Any data involving money movement or legal records MUST be Ledger tier.
- Last-write-wins MUST NOT be used for user-authored content without a product sign-off recorded in the design doc.
- Cross-service events MUST be published via transactional outbox or CDC; dual writes are prohibited.
- Services MUST emit
replication.lag.secondsandconsistency.violationmetrics (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#
| Signal | What 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)#
| Concept | Senior (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#
- Replicated Data Store — Quorums, leader election, replica lag, and failover data loss in depth
- Payment Processing — Ledger-grade consistency, idempotency, and reconciliation across providers
- Reservation Systems — Write skew, holds, and fail-closed booking
- Flash Sales & Ticketing — Escrow/allocation of inventory across regions
- Collaborative Editing — CRDTs vs OT and multi-region active-active editing
- Distributed Consensus — Raft/Paxos, fencing tokens, and linearizable reads
- Database Selection — Matching consistency guarantees to datastore choices
Related Foundations & Patterns: Consistent Hashing · Sharding & Partitioning · Dealing with Contention · Distributed Coordination
Related Technologies: PostgreSQL · Cassandra · DynamoDB · ZooKeeper & etcd