Technologies referenced in this case study: PostgreSQL · DynamoDB · Cassandra · Redis · Elasticsearch · Time Series DBs · Kafka
Related case studies: Database Sharding · Replicated Data Store · Distributed Caching · Search Indexing · Data Modeling · Database Indexing · Consistency Models · Build vs Buy
How to Use This Case Study#
Organized for interview use first, reference second. Read front-to-back once. Return to individual sections for targeted review.
| Mode | Time | What to Read |
|---|---|---|
| Quick Review | 15 min | Executive Summary → Interview Walkthrough → Fault Lines table → Drills 1–3 |
| Targeted Study | 1–2 hrs | Executive Summary → Walkthrough → Section 3 (Fault Lines) → Section 4 (Failure Modes) → the migration Deep Dive |
| Deep Dive | 3+ hrs | Everything, including the Principal Lens and appendices |
What is Database Selection? — Why interviewers pick this topic
Every system design interview contains a database decision, and most candidates make it in the first five minutes by reflex: "I'll use Cassandra because it scales," "MongoDB because the schema is flexible," "DynamoDB because it's serverless." Interviewers sometimes make it the whole question: "Which database would you use for X, and why?" or "Our Postgres is at 80% CPU. Should we move to NoSQL?"
Before vs After — the "we need to scale" scenario:
Technology-first choice:
Month 0: New messaging feature. Team picks a wide-column store "because it scales to millions of writes."
Month 2: Product asks for "unread count per conversation" and "search messages by sender." Neither is a
partition-key lookup. Team adds a second table per query, written by the application.
Month 5: The 4 denormalized tables drift: a bug writes 3 of 4. No transactions to fix it.
Month 8: Product asks for "edit message" with history. Another table. Read-repair storms at peak.
Month 12: Total load: 3K writes/s. A single Postgres node would have handled 10× this.
Month 14: Two engineers full-time on data consistency tooling. The on-call rotation hates the cluster.
Access-pattern-first choice:
Month 0: Team writes down the 8 queries, the write rate (3K/s peak, 10× in 2 years = 30K/s),
the correctness bar (per-conversation ordering, no lost messages) and team skills.
Picks Postgres: all 8 queries are indexed lookups; 30K/s fits a well-sized primary with batching.
Partitions by conversation_id from day one; states the migration trigger in the design doc:
"sustained > 50K writes/s or > 5 TB hot data → shard by conversation_id or move messages
to a wide-column store."
Month 12: 3K writes/s, one primary + 2 replicas. Features shipped in days, not weeks.
Month 30: Trigger hit for the messages table only. That one table migrates. The other 40 tables stay put.
Why interviewers reach for this question: it tests whether you choose tools from requirements or from reputation. It also tests whether you understand that a database is a 5–10 year commitment: operational expertise, migrations, backups, on-call and hiring all follow from the choice.
Mechanics Refresher: The Database Families
| Family | Examples | Sweet spot | Weakness |
|---|---|---|---|
| Relational (single-primary) | PostgreSQL, MySQL | Transactions, joins, ad-hoc queries, constraints; up to ~10s of TB, ~10–50K writes/s per primary | Write scale-out requires sharding; schema migrations at size need care |
| Distributed SQL (NewSQL) | Spanner, CockroachDB, YugabyteDB | SQL + transactions + horizontal scale | Higher write latency (consensus: ~5–20ms in-region, 50–150ms cross-region); cost; fewer operators know it |
| Key-value / document (managed) | DynamoDB, MongoDB | Known access patterns by key; predictable single-digit ms at any scale | Queries beyond the key design are expensive or impossible; hot partitions |
| Wide-column | Cassandra, ScyllaDB, Bigtable/HBase | Massive write throughput, time-ordered data per partition, multi-DC | No joins, limited transactions, query-first modeling, tombstones, repair ops |
| In-memory | Redis, Memcached | Sub-ms reads, counters, leaderboards, caches, ephemeral state | Memory cost ($/GB ~10–30× disk); durability is optional and weaker |
| Search | Elasticsearch, OpenSearch | Full-text, faceting, relevance | Not a system of record; eventual (refresh ~1s); reindexing is expensive |
| Time series | Prometheus, InfluxDB, TimescaleDB | Append-heavy metrics, downsampling, range queries | Cardinality explosions; limited updates |
| Analytical (OLAP) | ClickHouse, BigQuery, Snowflake | Scans and aggregations over billions of rows | Not for point lookups or high-rate small writes |
| Graph | Neo4j, custom (e.g., Facebook TAO over MySQL) | Multi-hop traversals | Niche; often better as a cache/index over a relational store |
For most production systems: PostgreSQL as the system of record, with purpose-built stores added as derived views (search, cache, analytics) fed by change data capture. Deviate from that default only with a written access pattern that Postgres can't serve at the projected 2–3 year scale.
Executive Summary
If you only read one section, read this. Everything else in the case study elaborates on the contrast below.
What This Interview Actually Tests#
Database selection is not a technology question. Every candidate can list SQL vs. NoSQL tradeoffs.
This is an access-pattern and commitment question that tests:
- Whether you derive the choice from queries, write rates, data size, and correctness bar instead of a store's reputation
- Whether you know the real scale envelope of the boring default, and so don't reach for a distributed store for 3K writes/s
- Whether you account for the cost of polyglot persistence: every additional store is another on-call rotation, backup regime, security review, and consistency boundary
- Whether you plan the migration path before you need it: the trigger, the mechanism, and the one-way doors
The key insight: the right database is the one whose failure modes and operational burden your team can own for five years while serving the access patterns you can write down today plus the growth you can defend. That's usually a relational database until a specific, measured access pattern proves otherwise.
The L5 → L6 → L7 Contrast — Start Here#
| Behavior | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| First move | Names a database ("Cassandra, because it scales") | Writes down access patterns, write/read rates, data size in 2–3 years, correctness bar, then chooses | Asks what the org already runs, operates well, and has on-call for; treats each new engine as a platform commitment |
| SQL vs NoSQL | "NoSQL for scale, SQL for consistency" | "Postgres handles ~10–50K writes/s and multi-TB on one primary. We're at 3K/s. NoSQL buys scale we don't need and costs queries we do need." | Sets the org default and the evidence bar for exceptions |
| Polyglot | One store per use case: search, cache, graph, time series | One system of record; other stores are derived via CDC, rebuildable, and explicitly eventually consistent | Caps the number of supported engines (e.g., 4–6) and funds platform teams for each |
| Consistency | "Eventual consistency is fine" | Names which operations need transactions/invariants and keeps them in one store | Defines data classes (money, identity, content, telemetry) with required guarantees org-wide |
| Migration | "We'll migrate if needed" | States the trigger metric and the migration mechanism (dual-write, backfill, shadow reads, cutover) up front | Budgets migrations as multi-quarter programs; tracks one-way doors; plans deprecations |
| Ownership | "The DBA team" | Service team owns schema and queries; platform owns engine ops, backups, upgrades; named owner per store | Decides which engines are paved roads with platform support vs. "you build it, you run it" |
Why "first move" separates levels
L5: Chooses from reputation. The choice might even be right, but the interviewer can't tell, because it isn't connected to a requirement. The first follow-up ("how would you query messages by sender?") exposes it.
L6: "Before picking a store, let me list the access patterns. Get conversation by ID, list messages in a conversation newest-first with pagination, unread count per user, search by text, and edit a message with history. Peak 5K writes/s growing to 30K in 2 years, 2 TB hot. Every one of these except full-text search is an indexed lookup in Postgres, and the scale fits. Search goes to a derived index."
L7: "We already operate Postgres and DynamoDB with platform teams behind both. A third engine needs a platform owner, on-call, backup/restore drills, security review, and hiring. That's ~2–4 engineers of ongoing cost. The bar for adding it is a workload neither existing engine can serve at the 3-year projection."
Why "polyglot" separates levels
L5: Adds a store per feature: Mongo for profiles, Redis for sessions, Elasticsearch for search, Neo4j for friends, Cassandra for feeds. Five stores, five failure modes, no transactional boundary between them. The application dual-writes to keep them consistent, and they drift.
L6: Separates the system of record from derived views. Writes go to one transactional store. Search, cache, and analytics are fed asynchronously through CDC and can be rebuilt from the source. The consistency contract of every derived view is stated: "search is up to ~2s stale."
L7: Sees polyglot as an org cost curve. Each engine needs expertise that doesn't transfer, so a company with 12 engines has 12 thin skill pools and no one who can debug any of them at 3am. Sets a supported-engine list and a retirement plan for strays.
Why "migration" separates levels
L5: Treats the choice as permanent or migration as trivial: "we'll just move to Cassandra later."
L6: Knows a live migration of a system of record takes 6–18 months: dual-write, backfill, verify, shadow-read, cut over reads, cut over writes, decommission. So the design names the trigger (a measured metric) and keeps the data access layer narrow, so a migration touches one module instead of 400 call sites.
L7: Plans the portfolio: which stores are on the way out, which migrations get funded this year, and which decisions are one-way doors (a sharding key, a vendor-proprietary data model) that need extra review.
The Staff Positions#
| Position | Rationale |
|---|---|
| Access patterns first, engine second | The engine is derived from the queries, rates, and correctness bar, never the reverse |
| PostgreSQL is the default system of record | Transactions, constraints, flexible queries, deep operator pool, and it handles more scale than most products reach |
| One system of record per entity | Invariants live in one transactional store. Everything else is a derived, rebuildable view |
| Derived stores via CDC, not application dual-writes | Dual-writes drift on partial failure. CDC from the log is ordered and replayable |
| Every new engine pays an operational tax | ~2–4 engineers of ongoing platform, on-call, backup and security effort per engine at a mid-sized company |
| Write the migration trigger into the design | "Sustained > 50K writes/s or > 5 TB hot → shard." It keeps the choice honest and the migration planned |
| Managed over self-hosted, unless scale or cost proves otherwise | Managed services cost 1.5–3× the raw infrastructure and save the team you'd need to run it |
The Three Intents#
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| Greenfield product | Access patterns will change monthly; speed of iteration | Relational default, narrow data-access layer, migration trigger documented | Premature distributed store that can't serve new queries | Transactions for anything involving money, identity, inventory |
| Known extreme workload | A specific access pattern at a scale the default can't serve (e.g., 500K writes/s of time-ordered events) | Purpose-built store for that workload only, modeled query-first | Using it for queries it wasn't modeled for; ops burden | Stated per workload (often per-partition ordering, eventual elsewhere) |
| Migrating off a store that hit its limits | Live system, zero downtime, data correctness during transition | Dual-write/CDC + backfill + verification + shadow reads + staged cutover | Silent divergence between old and new; a cutover you can't roll back | Byte-for-byte verified equivalence before cutover |
🎯 Staff Move: "I'll treat this as a greenfield choice with a known growth curve. I'll default to Postgres and only deviate for a specific access pattern I can write down that it can't serve at our 3-year projection. If one appears, only that workload gets a different store, fed from the system of record."
The Five Fault Lines#
| # | Fault Line | The Tension |
|---|---|---|
| 1 | General-Purpose vs Purpose-Built | Relational flexibility (any query, any time) vs. a specialized store's scale envelope (one query shape at extreme scale) |
| 2 | Transactions vs Horizontal Scale | Single-node ACID (simple invariants) vs. partitioned stores (scale, but invariants across partitions are your problem) vs. distributed SQL (both, at a latency and cost premium) |
| 3 | Polyglot vs Standardized | Best tool per workload vs. a small set of engines the org can operate deeply |
| 4 | Managed vs Self-Hosted (and Lock-In) | Managed: fast, reliable, expensive, proprietary APIs. Self-hosted: cheaper at scale, needs a team, portable |
| 5 | Evolve in Place vs Migrate | Stretch the current store (read replicas, partitioning, caching, sharding) vs. move to a new engine and pay a multi-quarter migration |
In the Wild: Real Production Systems#
Why this section belongs here: real migrations are the best evidence that database choice is about access patterns and operations, not reputation.
Discord — MongoDB → Cassandra → ScyllaDB for Messages#
Discord publicly described storing messages first in MongoDB, then moving to Cassandra (around 2017) when the data outgrew a single replica set, modeling messages by (channel_id, time bucket) partitions. Years later, with trillions of messages, they migrated again to ScyllaDB, citing garbage-collection pauses, hot partitions and operational toil with the Cassandra cluster, and they added an intermediate data-service layer to coalesce hot reads.
Staff insight: the access pattern (read recent messages in a channel, append new ones) never changed. The engine changed twice as scale and operational costs changed. Choose for the pattern, expect to revisit the engine, and keep a data-access layer that makes that revisit survivable.
Uber — PostgreSQL → MySQL (and Schemaless)#
Uber's engineering blog (2016) described moving core storage from PostgreSQL to MySQL, citing their workload's write amplification, replication behavior and upgrade difficulty at their scale. They built Schemaless, a sharded, append-only layer on top of MySQL.
Staff insight: this case is often cited as "Postgres doesn't scale." The more accurate lesson is that at very high write volume, engine internals such as MVCC, index write amplification and replication format become first-order, and big companies often build a thin custom layer over a boring engine rather than adopting an exotic one.
Figma and Notion — Scaling Postgres Instead of Replacing It#
Both companies publicly described staying on PostgreSQL through hypergrowth. Notion sharded its Postgres (2021) into hundreds of logical shards across a smaller set of physical databases, keyed by workspace. Figma first split tables across databases (vertical partitioning) and then built horizontal sharding on top of Postgres.
Staff insight: these are the counterexamples to "we'll outgrow SQL." When your data has a natural tenant key (workspace, file), sharding the relational database keeps transactions within a tenant and lets you keep the engine the team already knows. See Database Sharding.
What Interviewers Probe#
| After You Say... | They Will Ask... | (What They're Evaluating) |
|---|---|---|
| "Cassandra because it scales" | "What's the write rate? Would Postgres handle it?" | Do you know the default's real envelope? |
| "MongoDB for schema flexibility" | "What happens when half your documents have the old shape?" | Schema-on-read still has a schema, just an unmanaged one |
| "DynamoDB" | "Product wants a new query by a non-key attribute. What do you do?" | Query-first modeling cost; GSIs; lock-in |
| "Add Elasticsearch for search" | "How does it stay consistent with the database?" | Derived-store thinking; CDC vs. dual-write |
| "We'll migrate later" | "Walk me through migrating 20 TB live with zero downtime." | Migration mechanics and verification |
| "Postgres for everything" | "At what point does that stop being true? What's the trigger?" | Whether your default has known limits |
System Architecture Overview#
Reading the diagram: there's exactly one place each invariant lives: orders, payments and users in Postgres, inside transactions. The purpose-built store exists for one written-down workload and is modeled for that workload's queries. Everything else (search, cache, analytics) is a derived view fed by CDC from the log. It can be rebuilt from Kafka or from a snapshot, and its staleness is stated. The governance box matters as much as the data boxes: every engine on this diagram has a platform owner, a restore drill, and a documented migration trigger.
Quick-Reference: The 30-Second Cheat Sheet#
| Topic | The L5 Answer | The L6 Answer — Say This |
|---|---|---|
| Default | "Depends on the use case" | "Postgres as system of record unless a written access pattern needs otherwise at our 3-year scale." |
| Scale | "NoSQL scales better" | "One well-sized Postgres primary does ~10–50K writes/s and multi-TB. Replicas scale reads. Sharding by tenant scales writes." |
| Search / analytics | "Add Elasticsearch / a warehouse" | "Derived views via CDC: rebuildable, with a stated staleness. Never dual-write from the app." |
| Transactions | "Use eventual consistency" | "Money, inventory and identity live in one transactional store. Cross-store invariants are sagas with reconciliation." |
| New engine | "Best tool for the job" | "Best tool the org can operate. Each engine costs ~2–4 engineers ongoing. What's the evidence it's needed?" |
| Migration | "Migrate later if needed" | "Trigger is > 50K writes/s or > 5 TB hot. Mechanism: CDC + backfill + verify + shadow reads + staged cutover." |
Key Numbers Worth Memorizing#
| Metric | Value | Why It Matters |
|---|---|---|
| Postgres/MySQL single primary writes | ~10–50K simple writes/s (hardware and schema dependent) | Most products never exceed this |
| Postgres/MySQL single primary reads | 100K+ indexed point reads/s, more with replicas | Reads scale with replicas and caches before engines change |
| Comfortable single-node relational size | ~1–10 TB hot data | Above this, vacuum, backups, index builds and restores get slow |
| Replica lag (same region, healthy) | < 100ms typical | Read-your-writes needs routing care |
| DynamoDB per-partition throughput | ~3,000 RCU / 1,000 WCU | Hot keys throttle regardless of table capacity |
| DynamoDB max item size | 400 KB | Large blobs belong in object storage |
| Cassandra write throughput per node | ~10–30K writes/s (commodity) | Linear scale-out for append-heavy workloads |
| Cassandra partition size guidance | Keep < ~100 MB | Bucket time-series partitions by time |
| Redis single-thread throughput | ~100K+ ops/s per instance | Memory, not ops, is usually the limit |
| Elasticsearch refresh interval | 1s default | Search is near-real-time, not read-your-writes |
| Distributed SQL commit latency | ~5–20ms in-region, 50–150ms+ cross-region | The price of consensus |
| Live system-of-record migration | 6–18 months | Why the trigger and the data-access layer matter on day one |
| Operational cost of an additional engine | ~2–4 engineers ongoing (mid-sized org) | The hidden cost of polyglot persistence |
Interview Walkthrough
The most common mistake: the candidate names a database in minute two and spends the rest of the interview defending it. Treat the database choice as a conclusion, reached after the access patterns are on the board, and say out loud what would change your mind.
This walkthrough uses a concrete prompt: "Which database(s) would you use for a marketplace: listings, orders, messaging between buyers and sellers, search, and an activity feed?" The same shape applies when database selection is embedded in any other design.
Phase 1: Requirements & Framing (2–3 minutes)#
State the functional scope in one breath, then go straight to the numbers that decide engines:
"Before picking anything, I want five facts per workload: the queries, the read and write rates now and in 2–3 years, hot data size, the correctness bar, and whether the team already operates something that fits. Let me assume: 20M users, 5M listings, 200K orders/day, messages at ~2K writes/s peak, the feed at ~20K writes/s peak with fan-out, and search at ~3K queries/s."
Then set the correctness bar per workload. That's what separates the stores:
"Orders and payments need transactions and invariants: no double-sell of a unique item, and money balances. Listings need consistency for the owner and freshness within a few seconds for everyone else. Messages need per-conversation ordering and no loss. Search can be a couple of seconds stale. The feed can be minutes stale and lossy under load."
🎯 Staff Move: "I'm going to reach a database choice as a conclusion, not start with one. And I'll tell you what would change my mind: a measured trigger for each store."
Phase 2: Access Patterns & Data Model (3–4 minutes)#
This is the "entities and API" phase for this question. Write the access-pattern table:
| # | Access pattern | Rate (peak, 3-yr) | Consistency need | Shape |
|---|---|---|---|---|
| 1 | Create order (reserve listing, charge, record) | 50/s → 200/s | Transactional, multi-row | Multi-entity write |
| 2 | Get listing by ID | 20K/s | Owner: read-your-writes; others: seconds | Point read |
| 3 | Seller's listings, sorted, filtered | 2K/s | Seconds | Secondary index range |
| 4 | Full-text + faceted listing search | 3K/s → 10K/s | ~2s stale OK | Search |
| 5 | Messages in a conversation, newest first | 5K reads/s, 2K writes/s → 8K | Per-conversation order, no loss | Partitioned range |
| 6 | User activity feed | 50K reads/s, 20K writes/s → 150K | Minutes stale OK | Partitioned, time-ordered |
| 7 | Revenue by category, daily | Batch | Hours stale OK | Scan/aggregate |
"Now the choice falls out of the table. Patterns 1–3 are relational and transactional, and the rates are tiny, so they go in Postgres. Pattern 4 is a search engine, derived from Postgres by CDC. Pattern 5 fits Postgres today at 8K writes/s partitioned by conversation, with a trigger. Pattern 6 is the only one that grows past a single primary's comfort zone, at 150K writes/s in 3 years, so it's the one candidate for a purpose-built store. Pattern 7 goes to the warehouse via CDC."
🎯 Staff Move: Writing the rate and consistency columns next to the access patterns makes the choice defensible and makes deviations visible. Interviewers remember the candidate who built the table.
Phase 3: High-Level Architecture (≤5 minutes)#
Walk through it in 90 seconds:
- Postgres is the system of record for orders, listings, users and messages. The order flow is a single transaction:
SELECT ... FOR UPDATEon the listing, insert the order, update inventory status. - Messages live in Postgres, in a table partitioned by
conversation_idhash with an index on(conversation_id, created_at DESC). The trigger to move them out is "> 40K writes/s sustained or > 3 TB." - Feed goes to a wide-column store or DynamoDB, keyed
(user_id, time_bucket), because its 3-year write rate exceeds what a single primary handles comfortably and it needs no transactions. - Search, cache invalidation and analytics are derived from Postgres via log-based CDC into Kafka. Each consumer is rebuildable and has a stated lag.
- Every store has an owner, a restore drill, and a trigger in the design doc.
🎯 Staff Move: "This is two primary engines, not five. Search, cache and warehouse are views, not sources of truth. If any of them is lost, we rebuild it from Postgres and Kafka. That's the property that keeps polyglot persistence survivable."
Phase 4: Transition to Depth (1 minute)#
"The four interesting decisions are: the real limits of Postgres here and what the trigger is, the one workload that justifies a second engine, how derived stores stay consistent, and what the migration looks like the day a trigger fires. Which would you like first?"
If no preference, lead with the Postgres envelope and the trigger. It proves you know where the default stops.
Phase 5: Deep Dives (25–30 minutes)#
Deep dive A: How far does the default go? (6–7 min)
"A single Postgres primary on a large instance (say 64 vCPU, NVMe, 512 GB RAM) handles roughly 10–50K simple writes/s and 100K+ indexed point reads/s, depending on row width, index count and transaction size. The first things to hurt are usually not throughput. They're operational: vacuum on high-churn tables, index build time on multi-TB tables, backup and restore time, and connection count. So my ladder, in order, is: fix queries and indexes → add a connection pooler → add read replicas for stale-tolerant reads → cache hot reads → partition large tables (native partitioning by time or hash) → split databases by domain (vertical partitioning) → shard by tenant key → only then change engines for the specific workload."
Name the triggers: "> 70% sustained CPU on the primary at peak after tuning, > 5 TB hot, restore time > our RTO (say 1 hour), or replica lag > 1s at peak. Any of those opens a design review, not an automatic migration."
Deep dive B: The one workload that earns a second engine (6–7 min)
"The feed: 150K writes/s at year 3, append-mostly, read by partition (a user's recent items), no cross-partition transactions, tolerant of minutes of staleness. That's exactly the shape a wide-column store or DynamoDB is designed for. Model query-first: partition key (user_id, month) so partitions stay under ~100 MB, clustering key created_at DESC. Writes are idempotent by (user_id, item_id). The feed is also rebuildable from the event log, so it's not a system of record. Losing it is an availability incident, not a data-loss incident."
"Managed or self-hosted? At 150K writes/s, DynamoDB on-demand would cost real money. I'd model provisioned capacity with autoscaling against a self-hosted Cassandra/ScyllaDB cluster plus the 2–3 engineers to run it. At this size I'd pick managed. Self-hosting starts to win at several times this scale, if the org already has the expertise."
Deep dive C: Keeping derived stores consistent (6–7 min)
"The anti-pattern is the application writing to Postgres and then to Elasticsearch. If the second write fails, they diverge silently, and retries reorder updates. Instead, log-based CDC reads Postgres's WAL (Debezium-style) into Kafka keyed by entity ID, which preserves per-entity order. Consumers apply updates idempotently with a version check (if incoming.version > stored.version). Rebuild is a snapshot plus replay. Staleness is measured: cdc.lag_seconds per consumer, alert at > 30s. For read-your-writes on search, the listing owner's own view reads from Postgres, not the index."
Deep dive D: The migration when a trigger fires (5–6 min)
"Say messages hit 40K writes/s and we move them to a wide-column store. Steps: (1) the data-access layer already hides the store, so only it changes. (2) CDC from Postgres feeds the new store continuously. (3) Backfill history in the background, idempotently. (4) Verify with row counts plus checksums per conversation over sampled and then full ranges. (5) Shadow reads: serve from Postgres, also read from the new store and compare, and log mismatches. (6) Cut over reads per cohort: 1% → 10% → 100%. (7) Flip writes: the new store becomes primary, with reverse CDC or dual-write back to Postgres for a rollback window of 2–4 weeks. (8) Decommission. Realistically that's 2–3 quarters for a system of record."
Phase 6: Wrap-Up (2–3 minutes)#
"Two primary engines: Postgres for everything with invariants and for most workloads by default, and one purpose-built store for the feed because its write rate is the only one that outgrows a single primary. Search, cache and analytics are derived from the log and rebuildable. Every store has a written trigger that would make us revisit it, and a data-access layer that keeps the migration to one module."
The org close:
"The cost I'd flag isn't the license or the instances. It's that every engine needs people who can restore it at 3am. I'd rather run two engines deeply than five badly."
🎯 Staff Move: End on the operational commitment and the trigger. It shows the choice is a managed decision with an exit, not a bet.
Common Timing Mistakes#
| Mistake | L5 Does This | L6 Does This Instead |
|---|---|---|
| Picks engine in minute 2 | "I'll use Cassandra" | Builds the access-pattern table first; the engine is the conclusion |
| SQL vs. NoSQL lecture | 5 minutes on CAP and BASE | One sentence per engine family, tied to a pattern in the table |
| One store per feature | Mongo + Redis + ES + Neo4j + Cassandra | One system of record + derived views + at most one justified specialist |
| No numbers | "It needs to scale" | "2K writes/s now, 8K in 3 years, Postgres envelope 10–50K" |
| No exit | Treats the choice as permanent | States the trigger and the migration mechanism |
| No operations | Ignores backups, restores, on-call | Names owner, restore drill, and per-engine headcount cost |
1. The Staff Lens#
1.1 Why This Problem Exists in Staff Interviews#
Database choices are among the longest-lived decisions in engineering. Code gets rewritten in a year. Data and the engine holding it tend to outlive the team that chose them. A wrong choice doesn't fail on day one. It fails in year two, as a feature that can't be built, a consistency bug that can't be fixed without transactions, or an on-call rotation nobody wants. Interviewers use this topic to test whether a candidate's instinct is requirements-driven and operationally honest or reputation-driven and optimistic.
It's also the most common place where Senior candidates over-engineer. A candidate who reaches for a distributed store at 3K writes/s signals that they don't know what the boring default can do, and a Staff engineer's job is often to stop exactly that decision.
1.2 The L5 vs L6 vs L7 Contrast — Visual#
1.3 The Staff Question That Cuts Through Everything#
"Which access pattern can't the store we already run serve at the scale we'll have in three years, and who will operate the thing we add?"
- If no pattern fails the test, you don't add an engine. You tune, replicate, partition or shard what you have.
- If one pattern fails, you add a store for that pattern only, modeled for it and fed from the system of record where possible.
- If you can't name the operator, whether a platform team, a managed service or an on-call rotation, you can't add the engine regardless of fit.
2. Problem Framing & Intent#
2.1 The Three Intents — Explained#
Greenfield product. You don't know next quarter's queries. Flexibility dominates: ad-hoc queries, joins, constraints, transactions and schema migrations that take minutes. The relational default wins almost always. Discipline is still required: put a narrow data-access layer in front of it, pick a tenant key early (even unsharded), and write down the trigger.
Known extreme workload. You can write the single query shape and it's big: 500K event writes/s, 10B-row time series, sub-ms counters, full-text relevance. Here a purpose-built store earns its place, and the modeling is query-first: design the partition key for the read, accept that other queries will be expensive, and keep the workload's data derived or separable.
Migrating off a store that hit its limits. The hardest intent. The data is live, the old store is failing slowly, and correctness during the transition is the whole job. The design is the migration mechanism: CDC, backfill, verification, shadow reads, staged cutover, rollback window. See Deep Dive 4 and Appendix C.
🎯 Staff Move: "Most 'which database' questions are really greenfield questions with a growth curve. So I'll start from the default and make the interviewer's scale numbers argue me out of it."
2.2 When NOT to Add a New Database#
| Situation | Why a new engine is wrong | Do instead |
|---|---|---|
| Write rate < ~10K/s and data < ~2 TB | The relational default handles it with headroom | Index, pool, replicate |
| "We need flexible schema" | JSONB columns give flexibility inside a transactional store | Postgres JSONB + GIN index for the flexible part |
| "We need search" for simple prefix/filter queries | Trigram/full-text indexes handle modest search | Postgres full-text until relevance tuning or scale demands a search engine |
| "We need a graph" for 1–2 hop queries | Recursive CTEs and adjacency tables handle shallow traversals | Adjacency table + cache; graph engine only for deep traversal workloads |
| "We need a queue" | A database table as a queue works up to ~1–5K jobs/s with SKIP LOCKED | Postgres queue table until throughput or fan-out needs Kafka |
| Nobody will own it | Unowned engines become the incident nobody can debug | Don't add it |
🎯 Staff Move: "Postgres can do a surprising amount of what people add engines for: JSON, full-text, queues, geospatial with PostGIS, time partitioning. None of it is best-in-class. All of it avoids a new on-call rotation until scale proves otherwise."
2.3 What the Interviewer Leaves Underspecified#
| Underspecified | Why it matters | What I'd assume out loud |
|---|---|---|
| Read/write rates and growth | The single biggest input | Now and at 3 years, per workload |
| Data size and hot set | Memory vs. disk, restore time | Hot set in RAM for the primary store |
| Correctness per workload | Transactions vs. eventual | Money/inventory transactional; feeds eventual |
| Query flexibility needed | Relational vs. query-first | Greenfield → assume new queries monthly |
| Existing stack and skills | Operational cost of adding engines | Org runs Postgres and one managed KV store |
| Latency SLO | In-memory vs. disk, replicas | p99 < 50ms for reads on the hot path |
| Multi-region / residency | Engine support for geo-partitioning | Single region + DR first; residency later |
| Compliance (PII, retention, deletion) | Deletion across derived stores | Every derived store must support delete-by-user |
2.4 Precise Terminology#
| Term | Meaning | Why the precision matters |
|---|---|---|
| Access pattern | A query shape + its rate + its latency/consistency need | The unit of database selection |
| System of record (SoR) | The store whose data is authoritative for an entity | Invariants and transactions live here; everything else is derived |
| Derived store / view | A store populated from the SoR (search, cache, OLAP) | Must be rebuildable; staleness must be stated |
| CDC | Change data capture from the database log (WAL/binlog) | Ordered, replayable propagation without app dual-writes |
| Query-first modeling | Designing tables/partitions around specific reads (NoSQL) | New queries need new tables or indexes, a real cost |
| Hot partition | A partition key receiving disproportionate traffic | Caps throughput regardless of cluster size |
| Scale envelope | The range of rates/sizes where an engine is comfortable to operate | Different from "maximum benchmark number" |
| Trigger | A measured condition that opens a re-evaluation | Keeps the choice honest and migrations planned |
| Polyglot persistence | Using multiple engines in one system | Each adds an operational tax and a consistency boundary |
| One-way door | A decision that is very expensive to reverse | Shard keys, proprietary data models, and the SoR engine |
3. The Five Fault Lines#
3.1 Fault Line 1: General-Purpose vs Purpose-Built#
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Relational for everything | Any query, transactions, constraints, one engine to operate | Extreme write rates, huge time series, relevance search, sub-ms counters at scale | The team at the moment a workload crosses the envelope |
| Purpose-built per workload | Each workload at its best performance point | N engines, N skill sets, cross-store consistency, N backup regimes | Platform/on-call; product teams waiting on new queries |
| Relational SoR + purpose-built derived views | Flexibility and invariants in one place; specialists where proven | CDC pipeline to own; staleness to explain | Data platform team (CDC), modest |
| Relational SoR + one purpose-built primary for a proven workload | Scale where needed, simplicity elsewhere | That workload loses joins and transactions with the rest | The team owning that workload |
Staff default: relational SoR plus derived views. Add a purpose-built primary store only for a workload with a written access pattern whose 3-year projection exceeds the relational envelope, and keep that workload's invariants self-contained.
🧭 Principal Move: "The decision tree is the easy part. The hard part is making every team walk it the same way. I'd publish it as the paved road, with a required design review for any path that ends outside the supported-engine list."
3.2 Fault Line 2: Transactions vs Horizontal Scale#
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Single-primary ACID | Invariants are one BEGIN…COMMIT; simple reasoning | Write ceiling of one node; failover takes ~10–60s | Nobody, until the ceiling |
| Sharded relational by tenant | Transactions within a tenant; linear scale across tenants | Cross-tenant queries and transactions; resharding; the huge tenant | Platform team (shard routing, rebalancing) |
| Partitioned NoSQL (DynamoDB, Cassandra) | Near-unlimited throughput for key-based access | Invariants across partitions are application code; limited or expensive transactions | Application teams writing sagas, reconciliation and idempotency |
| Distributed SQL | SQL + serializable transactions + scale-out | ~5–20ms+ commit latency, cross-region worse; cost; smaller operator pool | Latency budget and the infra bill |
Staff default: keep invariants in one transactional scope. If data has a natural tenant key (workspace, merchant, account), shard relational by that key so transactions stay inside a shard. Choose distributed SQL when you genuinely need cross-entity transactions at scale-out sizes and can afford the latency, as in global inventory or ledgers spanning regions.
"The question I ask is: which invariants cross entities? 'Don't sell the same seat twice' is one row, fine anywhere. 'Debit A and credit B atomically' is two rows. If A and B can be on different partitions, I need a transaction across them, a saga with reconciliation, or a distributed SQL engine. I'd pick deliberately, not discover it in production."
🎯 Staff Move: "Eventual consistency isn't a feature you choose. It's a cost you accept. I want to know exactly which invariants I'm handing to application code before I accept it."
3.3 Fault Line 3: Polyglot vs Standardized#
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Teams choose freely | Best local fit; team autonomy | 10+ engines, shallow expertise, inconsistent backups, security gaps, stranded clusters when teams reorganize | Future on-call, security, and whoever inherits it |
| One engine for everything | Deep expertise, one toolchain | Workloads forced into a bad fit; heroic workarounds | Teams with extreme workloads |
| Supported catalog (4–6 engines) with exception process | Paved roads with platform support; exceptions possible with evidence | The catalog must evolve; exceptions need owners | Platform teams (one per engine), exception owners |
Staff default: a supported catalog, for example relational (Postgres), managed KV (DynamoDB), cache (Redis), search (Elasticsearch/OpenSearch), streaming log (Kafka), warehouse. Anything else needs a design review, a named owning team with on-call, backup/restore evidence, and a sunset plan.
🧭 Principal Move: "Each engine we support costs roughly 2–4 engineers a year in platform, upgrades, security patching and on-call, before any product work. Six engines is a 15–25 person investment. That's the budget conversation, and it's why 'best tool for the job' has to include 'the org can run it.'"
3.4 Fault Line 4: Managed vs Self-Hosted (and Lock-In)#
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Managed, portable engine (RDS/Aurora Postgres, managed Kafka) | Ops offloaded; exit path exists (it's still Postgres) | ~1.5–3× infra cost vs. self-run; less tuning control | Finance |
| Managed, proprietary (DynamoDB, Spanner, Cosmos DB) | Excellent scale and ops; deep cloud integration | Data model and API lock-in; exit means re-modeling and rewriting access code | Future migration team |
| Self-hosted open source | Cheapest at large scale; full control; portable | Needs a team: upgrades, repairs, backups, capacity, 3am pages | Platform team headcount |
Staff default: managed until the infrastructure savings of self-hosting exceed the fully loaded cost of the team to run it with margin. That usually means spend well into the millions per year on that engine and existing in-house expertise. Accept proprietary lock-in consciously for workloads where the managed service is materially better, and keep a data-access layer so the lock-in is in one module, not 400 call sites.
"DynamoDB lock-in is real, but the cost of the lock-in is the cost of the migration, and that is bounded by how many places touch the API. A repository layer makes lock-in a quarter's work instead of a year's."
3.5 Fault Line 5: Evolve in Place vs Migrate#
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Scale up / tune | Days of work; zero migration risk | Hard ceiling; big instances cost 2× per step | Finance, briefly |
| Replicas + caching | Reads scale ~5–10× | Stale reads; writes unchanged | App teams handling read-your-writes |
| Partitioning + domain split | Smaller tables, faster maintenance; spreads load across DBs | Cross-domain joins move to the app | Teams that owned cross-domain queries |
| Sharding the same engine | Writes scale; same skills, same SQL | Routing layer, resharding, cross-shard queries | Platform team, 2–4 quarters |
| Migrate to a new engine | Removes the ceiling for that workload | 6–18 months; dual-running cost; correctness risk; new ops skills | Everyone involved, and the roadmap |
Staff default: climb the ladder in order and stop at the first rung that clears the 3-year projection. Change engines only when the workload's shape (not just its size) is wrong for the current store, or when the operational cost of stretching exceeds the migration cost.
🎯 Staff Move: "A migration is a product with a roadmap, a team and an exit criterion. If I can't staff it that way, the right answer is another rung of the ladder on the current engine."
4. Failure Modes & Operational Reality#
4.1 The Wrong-Shape Store — Features You Can't Build#
Month 0: Orders stored in a wide-column store, partitioned by customer_id ("it scales").
Month 4: Finance needs "all orders for merchant X in March." Not a partition-key query.
Team adds orders_by_merchant table, written by the app alongside the first.
Month 6: Support needs "orders by status = stuck." A third table. A nightly full scan as backup.
Month 9: A partial failure writes orders but not orders_by_merchant. Merchant payouts are short.
Month 10: Reconciliation job built. Two engineers on data integrity for 6 months.
Month 16: Decision to migrate orders to Postgres. 3 quarters of work.
Detection: data.reconciliation_mismatch_total between denormalized tables; feature lead-time for "new query" requests; count of app-maintained denormalized copies per entity.
Blast radius: every downstream consumer of the entity. Here, merchant payouts.
Mitigation: introduce CDC from the primary table to derive secondary views instead of app dual-writes; run reconciliation jobs with alerting.
Prevention: access-pattern table in the design review; for entities with money or unknown future queries, a relational SoR by default.
Owner: the service team, with the design-review board owning the missed check.
4.2 Dual-Write Drift Between SoR and Derived Store#
The application writes Postgres, then Elasticsearch. A deploy restarts pods between the two writes, and 0.3% of listing updates never reach search. Retries reorder two updates to the same listing, so search shows the older price.
Detection: search.drift_ratio from periodic sampled comparison (1,000 random IDs every 10 min, comparing version numbers); user reports of "price in search ≠ price on page."
Blast radius: every derived view fed by dual-writes. Silent until users notice.
Mitigation: full reindex from the SoR snapshot (hours), then switch the feed to CDC.
Prevention: CDC from the WAL, ordered by key; versioned idempotent upserts; a continuous drift checker.
Owner: the search/data platform team owns the CDC pipeline; the listing team owns the version field.
4.3 Hot Partition in a Partitioned Store#
A celebrity's feed partition receives 40K writes/s. The DynamoDB partition caps at ~1,000 WCU, so throttling hits that key and adaptive capacity only partially helps. In Cassandra, the replicas owning that partition hit 100% CPU and coordinator timeouts spread to neighbors.
Detection: db.throttled_requests{partition}; top-N partition key sampling; p99 latency by key prefix.
Mitigation: write sharding (suffix the key with #0..#15, fan-in on read), or cache/aggregate before write; move pull-model fan-out to fan-out-on-read for celebrities.
Prevention: key design review asking "what is the maximum rate on a single key?"; load test with realistic key skew.
Owner: the service team (key design).
4.4 Relational Primary at the Ceiling — Slow Degradation#
CPU climbs 3% a month. Vacuum can't keep up on the 2B-row events table, bloat grows, and queries slow. A restore test takes 9 hours against a 1-hour RTO. Nobody pages, because nothing is down.
Detection: db.cpu_p95_at_peak trend; db.table_bloat_ratio; db.restore_drill_duration vs. RTO; db.replica_lag_p99.
Mitigation: archive or partition the events table by month (drop old partitions instead of deleting rows); move the events workload out; scale up as a bridge.
Prevention: capacity triggers with a 6-month runway alert; quarterly restore drills with the time recorded; per-table growth dashboards.
Owner: DB platform team (monitoring and drills) and service team (schema, archival).
4.5 Migration Divergence#
During a live migration, a backfill job and the CDC stream both write the same row. The backfill writes an older version after the CDC applied a newer one. 0.01% of rows are silently stale in the new store, and the cutover would have made them authoritative.
Detection: full verification pass with per-row version and checksum comparison after backfill completes and CDC catches up; shadow-read mismatch rate.
Mitigation: re-run verification; repair mismatches from the SoR; don't cut over until the mismatch rate is 0 over N days of shadow reads.
Prevention: version-guarded writes (only if incoming.version > current.version) for both backfill and CDC.
Owner: the migration team, with a single accountable lead.
4.6 Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Wrong-shape store | Denormalized copy count; reconciliation mismatches | All consumers of the entity | CDC-derived views; plan migration | Service team + design review |
| Dual-write drift | search.drift_ratio, sampled version compare | Every derived view | Reindex; switch to CDC | Data platform |
| Hot partition | db.throttled_requests{partition} | One key's users, sometimes the node | Write sharding; cache; fan-out-on-read | Service team |
| Primary at ceiling | CPU trend, bloat, restore time vs. RTO | Entire SoR | Partition/archive, split domain, scale up | DB platform + service |
| Migration divergence | Verification and shadow-read mismatches | Everything, if cut over | Version-guarded writes; repair; delay cutover | Migration lead |
| Unowned engine | Engine not in catalog; no restore drill on record | Whatever depends on it | Assign owner or migrate off | Platform leadership |
| CDC lag | cdc.lag_seconds > 30s | Search, cache, analytics freshness | Scale consumers; check slot/WAL retention | Data platform |
| Replication slot bloat (Postgres) | pg.replication_slot_retained_bytes | Primary disk fills → outage | Drop or advance stale slots; alert on retained WAL | DB platform |
🎯 Staff Move: "The CDC replication slot is the classic hidden coupling. A dead consumer makes Postgres retain WAL until the primary's disk fills. The derived store's outage becomes the system of record's outage unless someone owns that alert."
5. Evaluation Rubric#
5.1 Level-Based Signals#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Method | Picks by reputation, then justifies | Access-pattern table with rates, size, consistency; engine is the conclusion | Org-wide selection standard and decision tree; exceptions reviewed |
| Scale knowledge | "NoSQL scales" | Knows the relational envelope; climbs the ladder before changing engines | Knows the org's aggregate footprint and where it's heading |
| Consistency | "Eventual is fine" | Names invariants and keeps them in one transactional scope | Data classes with required guarantees by policy |
| Polyglot | Store per feature | SoR + CDC-derived views; one justified specialist | Supported-engine catalog with per-engine platform funding |
| Migration | "Migrate later" | Trigger + mechanism: CDC, backfill, verify, shadow, cutover, rollback | Migration portfolio, deprecations, one-way-door reviews |
| Operations | Mentions replicas | Owner, restore drills, CDC slot alerts, hot-key design review | Prices engines in headcount; decides build/buy/retire |
5.2 Strong Hire Signals#
| Signal | What It Sounds Like |
|---|---|
| Access patterns first | "Let me write the 7 queries with their rates before I pick anything." |
| Knows the default's envelope | "8K writes/s is comfortably inside one Postgres primary. I'll set a trigger at 40K." |
| Derived, not dual-written | "Search is a CDC-fed view. If it drifts, we rebuild it from the log." |
| Prices polyglot | "A third engine costs ~2–4 engineers a year. This workload doesn't justify it." |
| Plans the exit | "The data-access layer keeps a migration to one module. The trigger is in the design doc." |
5.3 Lean No-Hire Signals#
| Signal | Why It Misses the Bar |
|---|---|
| Engine named before any requirement | Reputation-driven; can't defend under follow-ups |
| Five engines for a modest workload | No sense of operational cost |
| App-level dual writes to keep stores in sync | Will drift on partial failure |
| "MongoDB because schemaless" with no migration story | Schema-on-read still has a schema, just an unmanaged one |
| No mention of backups, restores or ownership | Treats the DB as a library, not a system |
5.4 Common False Positives#
- Reciting CAP ≠ selecting a database. CAP describes partition behavior. It doesn't pick your engine.
- Knowing Cassandra's compaction strategies ≠ knowing when to use Cassandra.
- Benchmarks ≠ scale envelope. A vendor benchmark of 1M ops/s says little about your schema, your queries and your on-call team.
- "We'll use Spanner" ≠ solving consistency. It solves it at a latency and cost price that must be justified.
6. Interview Flow & Pivots#
6.1 Typical 45-Minute Shape#
| Phase | Time | Goal |
|---|---|---|
| Framing | 0–3 min | Workloads, rates, correctness per workload |
| Access patterns | 3–7 min | The table: query, rate now/3-yr, consistency, shape |
| Architecture | 7–12 min | SoR + specialist (if earned) + CDC-derived views |
| Deep dive 1 | 12–20 min | The default's envelope, the ladder, the triggers |
| Deep dive 2 | 20–28 min | The specialist workload: key design, hot partitions, managed vs self-hosted |
| Deep dive 3 | 28–36 min | Derived-store consistency: CDC, versions, drift detection |
| Deep dive 4 | 36–42 min | Migration mechanics |
| Wrap-up | 42–45 min | Operational ownership, engine count, triggers |
6.2 How Interviewers Pivot — And What They're Testing#
| Pivot | What They're Testing | Strong Response |
|---|---|---|
| "Now it's 100× the scale." | Whether the ladder is real | Which workload crosses its envelope first; shard by tenant or move that workload only |
| "Make it multi-region active-active." | Consistency vs. latency | Per-entity home region, or distributed SQL for global invariants, or CRDT-friendly data only |
| "The team only knows MongoDB." | Operational pragmatism | Skills are a real input; Mongo with transactions and schema validation may beat a Postgres nobody can run |
| "We need analytics on live data." | OLTP/OLAP separation | CDC to a warehouse; never heavy scans on the primary |
| "Cut the database bill by 40%." | Cost levers | Right-size, reserved capacity, archive cold data, remove unused indexes, move derived stores to cheaper tiers |
6.3 What to Deliberately Skip#
- Storage engine internals (LSM vs. B-tree) beyond one sentence tied to a read/write ratio.
- Vendor feature checklists.
- ORM discussions.
- Detailed index tuning. Link to Database Indexing if asked.
6.4 Follow-Up Questions to Expect#
- "Why not DynamoDB for everything?" New queries need new GSIs or tables, transactions are limited in scope, and hot keys and lock-in are costs. Good for known key-based patterns.
- "When would you pick MongoDB?" Document-shaped aggregates read and written as a unit, a team with Mongo expertise, and moderate cross-document invariants.
- "How do you handle read-your-writes with replicas?" Route the writer's reads to the primary for N seconds, or use LSN/GTID-based session consistency.
- "How do you do schema migrations on a 2 TB table?" Expand/contract: add nullable column, backfill in batches, dual-read, switch, drop. Use online schema-change tools.
- "How do you delete a user's data across all stores (GDPR)?" A deletion event via CDC/Kafka that every derived store consumes, with verification. Derived stores that can't delete aren't allowed.
- "What about the cache — is it a database?" A derived store with TTL and invalidation. Never the SoR. See Distributed Caching.
- "Who approves adding a new engine?" The platform architecture review, with evidence: access pattern, projection, owner, restore drill, exit plan.
7. Active Drills#
Drill 1: The Opening#
Prompt: "Which database would you use for a ride-sharing app?"
Staff Answer
"Let me split it into workloads first, because 'ride-sharing' is several databases' worth of access patterns. Trips and payments are transactional: a trip's state machine, fare, and charge must be consistent, so that's the relational system of record, sharded by city or rider if needed. Driver locations are ~1M updates every 4 seconds, which is ~250K writes/s of ephemeral data where only the latest value matters. That's an in-memory geo index (Redis with geo commands, or an in-memory service), not a durable database, with sampled history going to a log. Trip history for riders is read by rider and time, which fits the relational store at modest rates, or a wide-column store at very large scale. Analytics go to the warehouse via CDC. So: Postgres (or MySQL) as the SoR, an in-memory store for live locations, CDC to a warehouse. And I'd state triggers for when trips would need sharding."
Why this is L6:
- Decomposes into workloads with rates and consistency needs
- Recognizes ephemeral data that doesn't need a durable database at all
- Keeps invariants in one transactional store
What L7 adds:
- "Which of these engines do we already operate? Live location is the only workload where I'd accept a new operational burden."
- Plans residency: trips are per-city/country data with regulatory implications
❌ Common L5 Trap
"Cassandra for everything because it handles high write throughput and is highly available."
Why this misses: It's right for one workload (high write rate), wrong for trips and payments (need transactions), and wasteful for location (ephemeral, latest-value-only). The candidate optimized for the biggest number instead of the shapes.
Drill 2: The Envelope#
Prompt: "We're at 8K writes/s on Postgres and the team wants to move to Cassandra 'before it's too late.' What do you say?"
Staff Answer
"First, where does it hurt? CPU, I/O, vacuum, replica lag, connection count, restore time? 8K writes/s is well within a properly sized Postgres primary. If CPU is high, the likely culprits are missing indexes, too many indexes on hot tables, unbatched writes, or connection churn without a pooler. I'd take the ladder in order: query/index tuning, PgBouncer, batching, partitioning the largest table by time, moving stale-tolerant reads to replicas, splitting domains into separate databases. Each rung is weeks. A Cassandra migration is 2–3 quarters plus a new operational skill set, and we'd lose joins and transactions that the product almost certainly uses. I'd set a trigger instead: if sustained peak CPU exceeds 70% after tuning, or the growth projection crosses ~40K writes/s within 18 months, we open a design review on sharding or moving the specific hot workload."
Why this is L6:
- Diagnoses before prescribing
- Knows the envelope and the ladder
- Converts anxiety into a measurable trigger
What L7 adds:
- Quantifies the migration as
3–5 engineer-quarters plus a new on-call rotation, and compares it to the cost of a bigger instance ($10–20K/month) - Checks whether the org's paved road even includes Cassandra
Drill 3: Make It Concrete — Query-First Modeling#
Prompt: "You chose DynamoDB for chat messages. Show me the table design."
Staff Answer
"Access patterns first: (1) latest 50 messages in a conversation, paginated backwards; (2) append a message; (3) a user's conversations sorted by last activity; (4) unread count per conversation per user. Table messages: PK = conversation_id, SK = ts#message_id. That serves (1) with a descending query and (2) with a put. For long-lived busy conversations, add a time bucket to the PK (conversation_id#2026-09) to keep partitions and hot keys bounded. Pattern (3) gets a separate item type in a user_inbox table: PK = user_id, SK = last_activity#conversation_id, updated on each message. That's a second write, so it goes through a stream (DynamoDB Streams → consumer) rather than the app, to avoid drift. Unread count is a counter in the inbox item, updated by the same consumer and reset on read. Every new query product asks for later needs its own index or table. That's the cost I'm accepting in exchange for predictable latency at any scale."
Why this is L6:
- Access patterns drive keys
- Bounds partitions with time buckets
- Derives secondary views via streams, not app dual-writes
What L7 adds:
- Documents the "cost of a new query" so product understands the tradeoff up front
- Keeps a repository layer so the DynamoDB API is confined to one module (lock-in bounded)
Drill 4: Derived Store Consistency#
Prompt: "Search shows stale prices after sellers update listings. Fix it properly."
Staff Answer
"First find out whether it's lag or loss. Compare versions for a sample of listings between Postgres and the index. If the app dual-writes, it's probably loss on partial failure plus reordering on retries. The proper fix is log-based CDC: WAL into Kafka keyed by listing_id, so per-listing order is preserved, and a consumer that upserts with a version guard so older updates never overwrite newer ones. Then a drift checker that samples 1,000 IDs every 10 minutes and alerts if > 0.1% mismatch, and a full reindex path from a snapshot plus replay. Lag becomes a metric (cdc.lag_seconds, alert > 30s), and the seller's own listing page reads from Postgres, so they always see their own update."
Why this is L6:
- Distinguishes lag from loss
- Ordered, idempotent, versioned propagation
- Continuous verification plus rebuild path
What L7 adds:
- CDC as a platform service with SLOs, so each team doesn't build its own
- Replication-slot retention alerts owned by the DB platform, so a dead consumer can't fill the primary's disk
Drill 5: The Dependency Goes Down#
Prompt: "The primary database fails. Walk me through what happens."
Staff Answer
"With managed Postgres (Multi-AZ or a Patroni-style setup), failover to a synchronous standby takes ~30–60s. During that window, writes fail. The services should fail fast (short connection timeouts), return a clear error or queue non-critical writes, and not retry in a tight loop. Reads from replicas keep working, possibly stale. After failover, connection pools must reconnect to the new primary. A pooler behind a stable endpoint makes that transparent. The risks are a retry storm on reconnect, and async replicas that lost the last few hundred ms of commits, so I'd use synchronous replication to at least one standby for the SoR. Derived stores keep serving stale data. CDC resumes from the new primary's slot, which needs slot failover support or a re-snapshot. That last one surprises teams."
Why this is L6:
- Quantifies failover and data-loss window
- Covers client behavior, pools, retries and derived stores
- Catches the CDC slot failover gotcha
What L7 adds:
- Sets RPO/RTO per data class (payments: RPO 0; feed: minutes) and funds the replication topology accordingly
- Quarterly failover game days with the time recorded
Drill 6: The Hot Key#
Prompt: "One seller has 2M listings and 30% of all traffic. The store is sharded by seller_id. What now?"
Staff Answer
"Tenant-key sharding works until one tenant is a big share of the total. Options: move the whale to a dedicated shard (simple, buys time, but it's still one shard's ceiling), sub-shard the whale by (seller_id, listing_id hash) for that tenant only (transactions stay per-listing, which is fine for listing updates), cache their hot reads, and rate-limit their bulk writes through a queue. I'd do dedicated shard plus sub-sharding by listing, since their cross-listing transactions are rare. Longer term, the shard map should support per-tenant placement rules as a first-class feature."
Why this is L6:
- Recognizes the tenant-key failure mode
- Keeps transactional scope where it matters
- Graduated options with a choice
What L7 adds:
- Makes whale tenants a product/commercial concern: dedicated capacity in enterprise pricing
- Adds a tenant-size review to onboarding (see Database Sharding)
Drill 7: Multi-Tenant — Shared or Isolated?#
Prompt: "B2B SaaS with 5,000 customers. One database per customer or shared tables?"
Staff Answer
"Shared tables with tenant_id on every row and row-level security for the long tail. It's the cheapest, simplest migrations, and 5,000 schemas or databases means 5,000 migration runs and 5,000 connection pools. Enterprise customers who need isolation (contractual, residency, noisy-neighbor protection) get dedicated databases through the same code path: the tenant routing layer maps tenant_id → database. So the architecture is a pooled model plus a silo model behind one router. The tenant key is also the future shard key, so choosing it now makes sharding a routing change later."
Why this is L6:
- Picks a default and an exception path
- Counts the operational cost of per-tenant databases
- Sets up future sharding with the same key
What L7 adds:
- Prices silo tenants into the enterprise SKU
- Standardizes the tenant router as a platform component for all services
Drill 8: Build vs Buy — Managed or Self-Hosted?#
Prompt: "Our DynamoDB bill is $400K/month. Should we move to self-hosted Cassandra/ScyllaDB?"
Staff Answer
"Maybe, and the math decides. Self-hosted infra for the same workload might be $100–150K/month, but add a team: 4–6 engineers for a production cluster of that size with on-call, which is ~$150–250K/month fully loaded. Add a 3–4 quarter migration: dual-running cost plus engineers. Net savings might be ~$50–150K/month after the migration pays back in 1–2 years, and we take on durability and availability risk that AWS carried. Before that, I'd look for cheaper wins: provisioned capacity with autoscaling instead of on-demand (often 30–50% cheaper for steady load), reserved capacity, TTLs on data nobody reads, removing unused GSIs, and moving large attributes to S3. Those often cut 30–50% with no migration. If the bill is still growing fast after that, I'd run the build-vs-buy analysis formally."
Why this is L6:
- Includes team cost and migration cost, not just infra
- Exhausts cheaper levers first
- Names the risk transferred back to us
What L7 adds:
- Frames it as a 3-year TCO with a reversibility plan (see Build vs Buy)
- Considers negotiating an enterprise discount, often the cheapest lever of all
Drill 9: Changing the Engine Without an Outage#
Prompt: "Move the 12 TB orders table from MySQL to Postgres with zero downtime."
Staff Answer
"Phase 0: put every orders access behind one repository module; that alone can take a quarter. Phase 1: CDC from MySQL binlog → Kafka → Postgres, continuous. Phase 2: backfill history in chunks with version-guarded upserts, so backfill can never overwrite newer CDC data. Phase 3: verify with row counts per day partition, then per-row checksums, then continuous sampled comparison. Phase 4: shadow reads: every read goes to MySQL for the response and to Postgres for comparison, and mismatches are logged. Target: zero mismatches over 2 weeks. Phase 5: cut over reads by cohort, 1% → 10% → 50% → 100%. Phase 6: cut over writes. Postgres becomes primary, and CDC reverses (Postgres → MySQL) so we can roll back for 2–4 weeks. Phase 7: decommission. The cutover of writes is the one-way-ish step, so it happens during low traffic with a rehearsed runbook and a go/no-go checklist."
Why this is L6:
- Complete mechanism with verification gates
- Version-guarded writes prevent the backfill race
- Reverse replication makes rollback real
What L7 adds:
- Staffs it as a program with a single accountable lead and a budget
- Sets a decommission date up front; half-finished migrations running two engines forever are the most expensive outcome
Drill 10: Multi-Region#
Prompt: "We need the app in US and EU with low latency writes in both. What changes in the database layer?"
Staff Answer
"First question: which data needs global invariants? Most doesn't. User profiles, orders and messages belong to a user or tenant who has a home region. So home-region partitioning: each user's data lives in their region (which also satisfies residency), writes are local, and cross-region reads are rare and slower. For the few global invariants, like unique usernames or global inventory, either route those writes to one region (cross-region latency ~80–150ms on that path only) or use a distributed SQL engine with geo-partitioning. Active-active multi-master for general data means conflict resolution in application code, and I'd avoid it except for data that merges naturally (counters, sets)."
Why this is L6:
- Separates home-regioned data from global invariants
- Limits cross-region latency to specific operations
- Avoids multi-master unless the data is mergeable
What L7 adds:
- Aligns regions with legal boundaries and writes a residency standard
- Prices a distributed SQL engine vs. home-region routing over 3 years
8. Deep Dive Scenarios#
Deep Dive 1: Peak-Traffic Incident — The Primary Hits the Wall#
Context: During a seasonal peak, the Postgres primary for the core app is at 95% CPU, p99 query latency is 800ms (normally 15ms), and checkout is timing out. The on-call escalates to you.
Questions to Surface First:
- Which queries dominate load right now (
pg_stat_statements)? New or old? - Is it CPU from queries, lock contention, or connection overload?
- Are stale-tolerant reads hitting the primary instead of replicas?
- What changed recently: deploys, index drops, data growth, plan changes?
Typical L5 Approach: Scale the instance up (a 20–30 min failover risk mid-incident), add replicas, and propose moving to a "more scalable" database.
Staff Approach: Finds the top queries by total time. Typically one query with a regressed plan or a missing index dominates, or a new feature is doing N+1 queries on the primary. Kills or rate-limits the offending path with a feature flag, routes browse traffic to replicas, and sheds non-critical writes to a queue. After the incident: fix the index or plan, add query-level budgets, and model capacity against the next peak.
Principal Approach: Asks why one feature could consume the shared primary's capacity. Introduces per-service query budgets and statement timeouts per role, moves the SoR onto a domain-split topology so checkout's database is isolated from browse and recommendations, and adds DB capacity to the pre-peak readiness review.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Immediate (0–5 min) | pg_stat_statements top 10 by total time; check locks and connection count; feature-flag off the offending path. |
| Triage | Plan regression? Missing index? New feature? Read traffic on the primary? |
| Quick fix | Add or restore the index concurrently; route reads to replicas; set statement_timeout for non-critical roles. |
| Guardrails | Watch CPU, lock waits, replica lag; don't fail over mid-incident unless the primary is truly lost. |
| Post-mortem | Why wasn't the regression caught? Add query plan regression checks in CI against production-sized data. |
Metrics to Watch: db.cpu, pg.stat_statements.total_time_top10, pg.locks_waiting, db.connections_active, db.replica_lag_seconds, checkout.error_rate
Organizational Follow-up: Per-role statement timeouts; query review for new features touching the SoR; domain split roadmap.
Ownership Question: "Who can turn off a product feature to save the database?" Staff answer: the incident commander, pre-authorized by a runbook that lists which features are sheddable. Product agrees to that list before peak season.
Key Takeaway: "When the relational primary hits the wall, it's usually one query, not the engine. Find it before you reach for a new database."
What clears the Staff bar:
- Diagnoses via top queries before scaling
- Uses shedding and routing, not a risky failover
- Converts the incident into query governance
Deep Dive 2: Silent Failure — The Derived Store That Stopped#
Context: Customers report that newly created listings don't appear in search. It started ~6 hours ago. Postgres disk usage on the primary has grown 400 GB in those 6 hours and is at 88%.
Questions to Surface First:
- Is the CDC connector running? Is the replication slot advancing?
- How much WAL is the slot retaining?
- How long until the primary's disk fills?
Typical L5 Approach: Restart the search indexer; investigate search. Doesn't connect the disk growth.
Staff Approach: Recognizes the coupling immediately: the CDC connector died, its logical replication slot stopped advancing, and Postgres is retaining WAL for it. The search outage is about to become a system-of-record outage when the disk fills. Priority one is the primary: restart the connector if it can catch up quickly; otherwise drop the slot, which frees the disk and means the index must be rebuilt from a snapshot. Then rebuild search from snapshot plus fresh CDC.
Principal Approach: Treats derived-store pipelines as able to harm the SoR and sets a standard: every replication slot has an owner, a
max_slot_wal_keep_size(or equivalent) safety limit, and an alert on retained bytes. The org decides explicitly that the SoR's availability beats the derived store's completeness: slots are dropped automatically at a threshold, and derived stores must support rebuild.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Immediate | Estimate time-to-full. If < 1 hour, drop the slot and accept a reindex. Otherwise restart the connector and watch the slot advance. |
| Triage | Why did the connector die? OOM, schema change it couldn't parse, credentials? |
| Quick fix | Rebuild search from snapshot + new slot; communicate search staleness. |
| Guardrails | Set a WAL retention cap on slots; alert at 50 GB retained. |
| Post-mortem | Why wasn't cdc.lag_seconds paging? Why did a derived store's failure threaten the SoR? |
Metrics to Watch: pg.replication_slot_retained_bytes, pg.disk_used_pct, cdc.lag_seconds, search.newest_doc_age_seconds
Organizational Follow-up: Slot ownership registry; CDC platform with SLOs; schema change process that notifies CDC consumers.
Ownership Question: "Who is allowed to drop a replication slot on the primary?" Staff answer: the DB platform on-call, per a pre-approved runbook, because the SoR's availability outranks any derived store.
Key Takeaway: "Derived stores are supposed to be disposable. A replication slot without a retention cap makes them dangerous."
What clears the Staff bar:
- Connects search staleness to disk growth
- Prioritizes the SoR over the derived view
- Establishes guardrails as a standard
Deep Dive 3: Large Customer Onboarding — The Whale Tenant#
Context: A new enterprise customer will bring 3 TB of data and 25% of current write volume in a multi-tenant Postgres setup (pooled tables, tenant_id everywhere). Sales signed it; onboarding is in 6 weeks.
Questions to Surface First:
- What's their write pattern: steady, or bulk imports?
- Do they need isolation (contractual, compliance, residency)?
- Which tables will they dominate, and how do their indexes fit in memory?
- Can our tenant router place a tenant on a dedicated database today?
Typical L5 Approach: Scale up the shared database and import their data.
Staff Approach: Puts them in the silo model: a dedicated Postgres cluster behind the same tenant router, same schema, same code path. Imports via a throttled bulk path with verification. Their load and failures are isolated from 5,000 other tenants, and the shared database's working set stays in memory. If the router can't do per-tenant placement yet, that's the 6-week project.
Principal Approach: Makes tenant placement a product capability with pricing. Tenants above a size threshold (e.g., > 500 GB or > 5% of load) are automatically siloed, priced accordingly, and reviewed at deal time. Capacity is part of the sales process, not a surprise after signature.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Weeks 1–2 | Tenant router placement support; provision dedicated cluster; schema migrations applied to all clusters by the same pipeline. |
| Weeks 3–4 | Bulk import rehearsal on staging with their data volume; verify counts and checksums. |
| Week 5 | Production import, throttled; verification. |
| Week 6 | Go-live; monitor per-tenant latency and errors. |
Metrics to Watch: tenant.db_latency_p99{tenant}, shared_db.buffer_cache_hit_ratio, import.rows_per_sec, import.verification_mismatches
Organizational Follow-up: Tenant size review at deal desk; silo tier in pricing; migration pipeline that handles N clusters.
Ownership Question: "Who runs schema migrations across 1 shared + N silo databases?" Staff answer: the platform migration pipeline, not humans. Every cluster is migrated by the same automated process with per-cluster status.
Key Takeaway: "Design the tenant router so a whale is a placement decision, not a re-architecture."
What clears the Staff bar:
- Isolates rather than scales up
- Keeps one code path for pooled and siloed tenants
- Pushes capacity into the commercial process
Deep Dive 4: Post-Mortem — The NoSQL Choice That Couldn't Answer Questions#
Context: Two years ago, a team built the billing system on a wide-column store "for scale." Today billing has 400 writes/s, finance can't run the reports it needs, and a reconciliation bug overcharged 2,000 customers because two denormalized tables disagreed. You're asked to lead the post-mortem and the path forward.
Questions to Surface First:
- What were the access patterns and rates at design time? Were they written down?
- Who approved the choice, and was there a design review?
- How many denormalized copies of each entity exist, and who maintains them?
- What does finance need that can't be queried today?
Typical L5 Approach: Fix the reconciliation bug; add another table for finance's queries.
Staff Approach: Names the root cause as a selection error: billing is a correctness-critical, query-heavy, low-rate workload, the opposite of what the store is good at. Proposes migrating billing to Postgres with the full mechanism (CDC, backfill, verification, shadow reads, staged cutover), with a transactional model where invoice, line items and payments commit together. Keeps the wide-column store for the workload it suits, if any. Short term: reconciliation job with alerting and a manual review queue.
Principal Approach: Fixes the process that allowed it. Introduces a data-class policy (money, identity, inventory must live in a transactional SoR) and a design review gate for any new SoR engine with the access-pattern table as a required artifact. Also reviews other systems against the policy and funds the migrations in priority order.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Immediate | Refund affected customers; reconciliation job + alerting on mismatches. |
| Root cause | Selection mismatch: low-rate, correctness-critical, query-heavy workload on a query-first store. |
| Plan | Postgres billing schema; repository layer; CDC; backfill; verification; shadow reads. |
| Execute | 2–3 quarters, one accountable lead, finance as stakeholder. |
| Close | Decommission old tables; finance reports on the new store or a warehouse fed by CDC. |
Metrics to Watch: billing.reconciliation_mismatches, billing.invoice_total_vs_ledger_diff, migration.shadow_read_mismatch_rate
Organizational Follow-up: Data-class policy; design-review gate; audit of other SoRs.
Ownership Question: "Who owned the original decision?" Staff answer: the team made it, but the org lacked a gate. The post-mortem is blameless about the choice and specific about the missing review.
Key Takeaway: "Scale was never the constraint for billing. Correctness and queryability were. Choose for the constraint you actually have."
What clears the Staff bar:
- Identifies the selection error without blaming individuals
- Plans a full, verified migration
- Fixes the process with a data-class policy
Deep Dive 5: Multi-Region Expansion#
Context: The product runs on a single-region Postgres SoR. The company is launching in the EU with residency requirements and wants low latency for EU users within 9 months.
Questions to Surface First:
- Which data must stay in the EU: all user data, or specific classes?
- Which operations need global invariants (unique usernames, global catalog, cross-region payments)?
- Do users ever collaborate across regions?
Typical L5 Approach: Set up a cross-region replica of Postgres in the EU and route EU reads there; writes still go to the US.
Staff Approach: Home-region partitioning. Each user/tenant has a home region stored in a small global directory. EU users' data lives in an EU Postgres cluster with the same schema and code. The global directory and a few global invariants (username uniqueness) live in a small global store, which is either a single-region service with caching (cross-region write latency only at signup) or a distributed SQL engine. Derived stores (search, analytics) are per-region, with an aggregated, residency-compliant analytics layer.
Principal Approach: Writes the residency standard: data classes, allowed cross-region flows, audit logging, and a migration workflow for users changing region. Sequences the rollout: new EU users first, existing EU users migrated by cohort. Also chooses whether to invest in distributed SQL as a platform for future regions or keep the home-region pattern.
Staff Approach — Full Reasoning
| Phase | What to Do |
|---|---|
| Design | Home-region directory; per-region SoR clusters; identify global invariants. |
| Build | EU cluster; tenant router extended with region; per-region CDC and search. |
| Migrate | Existing EU users via the standard migration mechanism, cohort by cohort. |
| Verify | Residency audit: no EU personal data in US stores (derived stores included). |
| Operate | Per-region on-call; region-level restore drills. |
Metrics to Watch: directory.lookup_latency_p99, crossregion.write_latency_p99, residency.violations_total (must be 0), migration.users_moved
Organizational Follow-up: Residency standard; region as a first-class dimension in the data platform.
Ownership Question: "Who certifies residency compliance?" Staff answer: legal/privacy defines the rules; the data platform provides the audit tooling and signs off per release.
Key Takeaway: "Multi-region databases are mostly a data-classification exercise. Home-region most data, and pay for global consistency only where an invariant demands it."
What clears the Staff bar:
- Separates home-regioned data from global invariants
- Includes derived stores in residency
- Reuses the tenant router and migration mechanism
9. Level Expectations Summary#
After studying this case study, you should be able to:
- Build an access-pattern table (query, rate now and at 3 years, size, consistency, shape) and derive the engine from it
- State the realistic envelope of a relational primary and climb the scaling ladder before changing engines
- Keep invariants in one transactional scope, and name exactly which invariants you hand to application code when you don't
- Design derived stores via CDC with versioned idempotent updates, drift detection and a rebuild path
- Model a key-value/wide-column table query-first, with bounded partitions and hot-key mitigation
- Execute a zero-downtime migration: repository layer, CDC, backfill, verification, shadow reads, staged cutover, reverse replication
- Price polyglot persistence in headcount and argue managed vs. self-hosted with a TCO
The Bar for This Question#
Mid-level (L4): Knows SQL vs. NoSQL differences, picks a reasonable store for a single workload, designs a schema.
Senior (L5): Chooses competently per workload and knows each engine's strengths. But the choice is often made from reputation, polyglot costs are ignored, derived stores are kept in sync with app dual-writes, and migration is "later."
Staff+ (L6): Reaches the choice as a conclusion from written access patterns and rates, defaults to a relational SoR, adds a specialist only for a workload that crosses the envelope, derives the rest via CDC, and documents triggers and a migration mechanism. Names owners, restore drills and operational cost. The interviewer should learn something from the answer, whether that's the replication-slot coupling, the backfill-vs-CDC version race, or how little a feature needs to justify staying on Postgres.
10. Staff Insiders: Controversial Opinions#
10.1 "Postgres Is the Right Answer More Often Than Anyone Admits"#
| Evidence | Implication |
|---|---|
| Notion and Figma publicly scaled through hypergrowth on sharded Postgres | The ceiling is further than folklore says |
| Most products never exceed ~10K writes/s on their SoR | The distributed-store premium buys nothing |
| JSONB, full-text, partitioning, PostGIS and queue tables cover many "new engine" asks | Fewer engines, fewer rotations |
The Staff position: Postgres by default. Make the workload prove otherwise with numbers.
Why this matters in interviews: defending the boring choice with an envelope and a trigger is more senior than picking the exotic one.
10.2 "Schemaless Is a Lie — You Just Moved the Schema Into Every Reader"#
| Evidence | Implication |
|---|---|
| Documents written 2 years ago have old shapes | Every reader handles N versions |
| Validation happens in app code, inconsistently | Bad data enters silently |
| Migrations still happen: lazily, forever | Schema debt accumulates |
The Staff position: use schema-on-read where data is genuinely heterogeneous. Otherwise enforce schema at write, including in document stores (schema validation), or use JSONB inside a relational table for the flexible part only.
Why this matters in interviews: "flexible schema" is the weakest common justification for a document store.
10.3 "Most Polyglot Architectures Are Org Charts, Not Designs"#
| Evidence | Implication |
|---|---|
| Each team picked its favorite store | Engine count tracks team count |
| Nobody can restore half of them confidently | Operational risk hidden until an incident |
| Cross-store consistency is dual-write glue | Silent drift |
The Staff position: one SoR per entity, derived views via CDC, a supported-engine catalog, and a sunset plan for strays.
Why this matters in interviews: it shows you see the org cost of technical choices.
10.4 "The Database Migration Is Never the Hard Part — The Access Layer Is"#
| Evidence | Implication |
|---|---|
| Data movement is mechanical: CDC, backfill, verify | Tooling exists |
| Finding and changing 400 call sites across 30 services is not | Most of the timeline is here |
| A repository layer turns lock-in into a bounded task | Day-one decision, multi-year payoff |
The Staff position: a narrow data-access layer per domain from day one. It's the cheapest insurance in the system.
Why this matters in interviews: answering "how would you migrate?" with "first, there's one module to change" signals you've done one.
10.5 "Discord's Two Migrations Are a Success Story, Not a Cautionary Tale"#
| Evidence | Implication |
|---|---|
| The access pattern stayed constant; the engine changed as scale and ops costs changed | Selection is not forever |
| They adapted the layer in front of the store (a data service to coalesce hot reads) | The access layer made the moves survivable |
| Each move was driven by measured pain | Triggers, not fashion |
The Staff position: expect to revisit engine choices every few years at hypergrowth. Design so revisiting is affordable.
Why this matters in interviews: it reframes "choose right the first time" into "choose defensibly and keep the exit cheap."
11. The Principal Lens (L7)#
Why L7 Sees This Problem Differently#
A Staff engineer picks the right database for a system. A Principal engineer sees that the org's data engines are a portfolio of long-lived commitments. Each carries a platform team, a skill pool, a security surface, a backup and restore regime, a vendor relationship, and a migration liability. The unit of decision is the catalog: which engines the company supports, how deeply, for which data classes, and which are being retired. Individual selection decisions get easy once the catalog, the data-class policy and the review gate exist. Without them, 40 teams make 40 locally reasonable choices that add up to an unoperable estate.
The Org-Level Fault Line#
Team autonomy in data-store choice vs. a governed platform catalog. Autonomy lets teams move fast and fit tools to problems. It also produces engine sprawl, orphaned clusters after reorgs, inconsistent backups and a security team that can't audit what it can't see. A strict single-engine mandate avoids sprawl but forces bad fits and invites shadow IT.
The Principal position: a supported catalog of 4–6 engines, each with a funded platform team, paved-road tooling (provisioning, backups, CDC, observability, access control), and published guidance on when to use it. Off-catalog engines are allowed through a review that requires a named owning team with on-call, restore evidence, a cost estimate, and an exit plan. Every off-catalog engine is re-reviewed yearly.
🧭 Principal Move: "I'd rather say yes to a new engine with an owner and an exit plan than say no and find it running unowned in production 18 months later. The review isn't a gate to stop teams. It's a way to make the cost visible before it's committed."
Cost Model#
Assumptions: managed-service list prices; fully loaded engineer ~$300K/year; "ops" includes upgrades, patching, backups/restore drills, capacity planning and on-call.
| Scale | Engines in use | Infra $/month | Platform headcount | On-call | Typical failure |
|---|---|---|---|---|---|
| Startup (≤ 50 engineers) | 2–3 (managed Postgres, Redis, maybe search) | $10–40K | 0.5–1 FTE spread across teams | Shared rotation | Missing restore drills; single primary without tested failover |
| Growth (~300 engineers) | 5–8 if ungoverned; 4–5 with a catalog | $300K–1M | 8–15 (2–3 per supported engine) | One rotation per engine | Polyglot sprawl; dual-write drift; CDC slot incidents |
| Hyperscale (2,000+ engineers) | Catalog of 5–7 + a few justified exceptions | $5–20M | 60–120 across data platform teams | Per engine, per region | Migration debt; lock-in; multi-region residency |
Reading it: at growth scale, each unnecessary engine costs ~$0.6–1M/year in people before infra. Consolidating two stray engines often funds a whole platform team.
The 3-Year Evolution Path#
One-Way Doors vs Two-Way Doors#
| Decision | Door | Reversibility Cost |
|---|---|---|
| System-of-record engine for money/identity | One-way | 2–4 quarters of migration with correctness risk |
| Shard key / tenant key | One-way | Resharding touches every row and every query path |
| Proprietary data model (DynamoDB single-table design) | One-way-ish | Re-modeling plus rewriting the access layer. Bounded if a repository layer exists |
| Adding a derived store fed by CDC | Two-way | Drop it and rebuild elsewhere from the log |
| Cache layer | Two-way | Remove it and check the SoR can take the load |
| Managed vs. self-hosted for a portable engine | Two-way | Standard replication tooling moves it |
| Region/home-region assignment of data | One-way per user, two-way per policy | Moving users is a migration; changing policy for new users is cheap |
The Standard I'd Write#
RFC: Data Store Selection & Operation Standard v1
Scope: Any new data store, or any new system of record, in production.
MUST
- Every design doc includes an access-pattern table: query, rate now and at 3 years, hot data size, consistency requirement, latency SLO.
- Money, identity, inventory and entitlement data (Class A) live in a transactional system of record from the supported catalog.
- Each entity has exactly one system of record. Other copies are derived via the CDC platform, must be rebuildable, and must document staleness.
- No application dual-writes to keep stores consistent.
- Every store has a named owning team, on-call, a restore drill at least quarterly with recorded duration vs. RTO, and a documented migration trigger.
- Every replication slot / CDC consumer has an owner and a retention cap.
- Access to each store goes through a per-domain data-access layer.
SHOULD
- Prefer managed offerings of catalog engines unless a TCO review shows > 30% savings net of team cost.
- Tenant/owner key present on every row from day one.
Supported catalog v1: PostgreSQL (SoR), DynamoDB (high-scale key-value), Redis (cache/ephemeral), OpenSearch (search, derived), Kafka (log/CDC), warehouse (analytics, derived).
Exceptions: Architecture review with the access-pattern table, owner, cost estimate and exit plan. Re-reviewed annually.
Success metrics: off-catalog engines ≤ 3 org-wide; 100% of SoRs with a passing restore drill in the last quarter; zero dual-write drift incidents; migration triggers documented for 100% of SoRs.
What I'd Tell the VP#
"Our biggest data risk isn't that we picked a database that can't scale. It's that we run eleven different kinds of databases, several with no clear owner and no tested restore. Every extra kind costs us roughly two to four engineers a year just to keep safe. I want to standardize on five supported types, give each one a platform team, and move the strays onto them over the next 18 months. That should free up around 8–10 engineers and cut the chance that a data-loss incident hits a system nobody knows how to restore. For new projects, teams will write down their data needs before choosing a store, which is cheap and prevents the expensive rewrites we did on billing last year."
Principal Interview Signals#
| Signal | What It Sounds Like |
|---|---|
| Portfolio view | "The question isn't which database for this service. It's whether this service justifies a sixth engine in our catalog." |
| Prices engines in people | "Each engine is 2–4 engineers a year. Two strays cost us a platform team." |
| Data-class policy | "Money and identity live in a transactional SoR. That's a policy, not a per-team debate." |
| Plans retirement | "We deprecate the stray document store over 3 quarters, with its three consumers migrating first." |
| Knows the one-way doors | "The tenant key gets an architecture review. The cache layer doesn't." |
Staff answers that L7 interviewers find insufficient:
- "I'd pick Postgres here and add Elasticsearch for search." Right for the system, but silent about the org's catalog, platform ownership and the cost of the search cluster.
- "We'll migrate when we hit the trigger." That doesn't budget or staff the migration, or say what else it displaces on the roadmap.
- "Managed is easier." That's true, but without a TCO it isn't a decision, and it ignores lock-in on proprietary data models.
Appendices
Appendix A: Engine Families in Depth#
A.1 Relational (Postgres/MySQL)#
Right when: multi-entity invariants, evolving queries, moderate to high scale within one primary or a tenant-sharded fleet. Wrong when: the workload is a single extremely high-rate append pattern with no joins, or relevance search.
-- Order placement: the invariant lives in one transaction
BEGIN;
SELECT status FROM listings WHERE id = :listing FOR UPDATE; -- must be 'available'
INSERT INTO orders (id, listing_id, buyer_id, amount, status) VALUES (...,'pending');
UPDATE listings SET status = 'sold' WHERE id = :listing;
INSERT INTO outbox (event_type, payload) VALUES ('order_created', ...); -- for CDC consumers
COMMIT;
The transactional outbox gives downstream systems an ordered event without dual-writes.
A.2 Key-Value / Document (DynamoDB, MongoDB)#
Right when: access is by known keys, latency must stay flat at any scale, items are self-contained aggregates. Wrong when: ad-hoc queries, multi-entity transactions at scale, or unknown future access patterns.
Table: user_inbox
PK = user_id SK = last_activity#conversation_id
attrs: unread_count, last_message_preview
Query: PK = :user ORDER BY SK DESC LIMIT 20
A.3 Wide-Column (Cassandra/ScyllaDB/Bigtable)#
Right when: very high write throughput, time-ordered data per partition, multi-DC replication. Wrong when: you need joins, secondary queries or strong multi-row transactions. Also wrong for delete-heavy workloads (tombstones).
CREATE TABLE feed (
user_id bigint, bucket text, -- bucket = 'YYYY-MM' bounds partition size
created_at timestamp, item_id uuid, payload blob,
PRIMARY KEY ((user_id, bucket), created_at, item_id)
) WITH CLUSTERING ORDER BY (created_at DESC);
A.4 Distributed SQL (Spanner/CockroachDB/YugabyteDB)#
Right when: you need cross-entity transactions and horizontal scale, possibly across regions. Global ledgers and inventory are examples. Wrong when: a single primary or tenant sharding suffices. You'd pay consensus latency and cost for nothing.
A.5 Where Each Family Sits#
Appendix B: The Access-Pattern Worksheet#
| Field | Question | Example |
|---|---|---|
| Query | Exact shape, including filters and sort | Messages in conversation, newest first, 50 per page |
| Rate now / 3-yr | Peak reads/s and writes/s | 5K / 20K reads; 2K / 8K writes |
| Data size now / 3-yr | Hot and total | 300 GB / 2 TB hot |
| Consistency | Transactional? Read-your-writes? Staleness allowed? | Per-conversation order; no loss |
| Latency SLO | p99 | 50ms |
| Invariants | What must never be violated, across which entities | Message IDs unique per conversation |
| Lifecycle | Retention, deletion, legal hold | 7 years; delete on user request |
| Residency | Region constraints | EU users' data in EU |
| Owner | Team and on-call | Messaging team |
| Trigger | What reopens the decision | > 40K writes/s or > 3 TB hot |
Appendix C: Migration Mechanics#
Rules:
- Both backfill and CDC use version-guarded upserts, which removes the older-overwrites-newer race.
- Verification runs after CDC has caught up past the backfill's end.
- Write cutover always has a rollback path (reverse replication) for a defined window.
- The old store has a decommission date set at the start.
Appendix D: Consistency Contracts for Derived Stores#
| Derived store | Feed | Staleness SLO | Read-your-writes strategy | Rebuild |
|---|---|---|---|---|
| Search | CDC → Kafka → indexer | p99 < 5s | Owner views read from SoR | Snapshot + replay, ~hours |
| Cache | CDC invalidation + TTL | < 1s after invalidation; TTL bound 5 min | Write-through for the writer's session | Cold start from SoR |
| Warehouse | CDC → batch/stream load | < 15 min | N/A | Reload from snapshots |
| Feed store | Events → fan-out consumers | < 1 min | Author's own post injected client-side | Replay from event log |
Appendix E: Observability#
E.1 Core Metrics#
# System of record
db.cpu_p95_at_peak, db.connections_active, db.replica_lag_seconds
pg.stat_statements.top_total_time, pg.table_bloat_ratio
db.restore_drill_duration_minutes (vs RTO)
db.growth_bytes_per_day → runway_days to trigger
# CDC and derived stores
cdc.lag_seconds{consumer}, pg.replication_slot_retained_bytes{slot}
derived.drift_ratio{store} # sampled version comparison
# Partitioned stores
db.throttled_requests{partition}, db.hot_key_topN
# Migrations
migration.shadow_read_mismatch_rate, migration.backfill_progress_pct
E.2 Critical Alerts#
| Alert | Threshold | Severity |
|---|---|---|
| Replication slot retained WAL | > 50 GB or > 20% free disk | Page |
cdc.lag_seconds | > 30s for 10 min | Page (search/cache), ticket (warehouse) |
derived.drift_ratio | > 0.1% | Page |
| Restore drill duration | > RTO | Ticket, escalated if two quarters in a row |
| Runway to trigger | < 6 months | Ticket to owning team + capacity review |
| Replica lag | > 5s for 5 min | Page if read-your-writes depends on replicas |
Appendix F: Scale Evolution#
| Stage | Works | Don't build yet |
|---|---|---|
| Early | One managed Postgres, repository layer, tenant key, backups tested | Sharding, specialist stores, CDC platform |
| Growth | Replicas, pooler, partitioning, CDC to search/warehouse, cache with invalidation | Distributed SQL, self-hosting |
| Scale | Tenant sharding, one specialist per proven workload, engine catalog | Custom storage engines |
| Global | Home-region partitioning, residency audit, managed-vs-self-hosted TCO reviews | Multi-master for non-mergeable data |
What you don't build on day one: a polyglot architecture. Do build on day one: the repository layer, the tenant key, tested backups and a written trigger.
Appendix G: Multi-Tenancy, Fairness & Cost#
- Pooled vs. silo: pooled tables with
tenant_id+ row-level security for the long tail. Dedicated databases for whales and compliance tenants, behind the same router. - Noisy neighbors: per-tenant statement timeouts and connection limits; per-tenant query budgets on shared primaries.
- Cost attribution: storage bytes, IOPS and query time per tenant, which feed pricing tiers.
- Deletion: per-tenant and per-user deletion events flow through CDC to every derived store; deletion completeness is verified, not assumed.