Databases & Storage
Choose SQL or NoSQL by access pattern, then add indexes, transactions, replication, caching and sharding, and know what each one costs
SPACED REPETITION Β· 16 practice questions
Make this lesson stick.
Try 3 questions now. No account needed. Sample answers aren't saved.
or sign in to practice all 16The database rarely breaks. One workload inside it does.
Suppose you launch a ride-hailing app on a single PostgreSQL server. Riders, drivers, trips, payments and live driver positions are all tables in one database, and in one city it works beautifully. Eighteen months later you are in 40 cities. These are the numbers you would state as assumptions in an interview:
- 1 million riders a day, 2 trips each: 2 million trips a day.
- Each trip writes 5 times: requested, matched, started, finished, paid.
- 500,000 drivers online at the evening peak, each sending a GPS position every 4 seconds.
- 50,000 riders watching the map at peak, which refreshes "cars near me" every 5 seconds.
- Riders open their trip history 4 times a day. Peak traffic is 3Γ the daily average.
This is what the one server is being asked to do at peak:
| Workload | Peak rate | Data it touches | Useful for how long? |
|---|---|---|---|
| Driver GPS updates | 125,000 writes/s | 500,000 current positions, about 50 MB | About 10 seconds |
| "Cars near me" searches | 10,000 reads/s | The same 50 MB | Seconds |
| Trip lifecycle writes | about 350 writes/s | 2.19 billion rows over 3 years, about 1.1 TB | Years (it is money) |
| Trip-history reads | about 140 reads/s | One rider's last 20 trips | Years |
Predict first: which row breaks the single-server design first?
Check your answer
The GPS updates. 125,000 writes a second is 360 times the trip write rate, and more than the few thousand to low tens of thousands of writes a second you would assume one primary handles (Move 1). Even hardware that could take it would pay the full price of durable storage on every one. In PostgreSQL an UPDATE writes a new version of the row, logs it to the write-ahead log, ships it to every replica and leaves a dead row version for vacuum to clean up. All that care protects data that is worthless ten seconds later and that the next ping will replace anyway.
The trip data, by contrast, is boring: a few hundred writes a second and about a terabyte over three years fit on one well-provisioned primary with its standbys. So the fix is not "switch to NoSQL". It is "move the one dataset whose access pattern does not fit". The whole set of current positions is 50 MB, which fits in the RAM of any server, although Move 1 shows that its command rate needs more than one Redis node.
Say it in the interview: "125,000 updates a second is beyond what I'd assume one Postgres primary handles. And even if big hardware could take it, I wouldn't spend its write capacity, log volume and vacuum on data that is stale in ten seconds and rebuilt by the next ping."
That is the whole lesson in one sentence: a storage design is a set of datasets, and each one is placed according to how it is read, how it is written, and what it costs to lose it or serve it stale. Six moves cover almost every storage decision in an interview. Each starts with a pressure in the requirements:
| # | Pressure in the requirements | Move | What you give up |
|---|---|---|---|
| 1 | Datasets with very different shapes, rates or lifetimes | Place each dataset by access pattern (relational by default) | Another system to run and keep in sync |
| 2 | A frequent query filters or sorts and ends up scanning | Add the index that matches the query | Slower writes, more storage |
| 3 | Writes must succeed together, or requests race for one row | Transactions, conditional writes, the right isolation | Contention and retries, inside one database only |
| 4 | A node can die; reads outgrow one machine | Replicate | Lag, stale reads, lost writes on async failover |
| 5 | The same data is read far more often than it changes | Cache | Staleness, invalidation bugs, stampedes |
| 6 | Data or writes outgrow one primary | Partition (shard) by the key your queries use | Cross-shard queries, hot keys, resharding |
This lesson assumes you know where load balancers, stateless services and object storage sit (Core Building Blocks) and can turn users into requests per second (Foundations of System Design). CAP, PACELC and consistency models are defined precisely in Key Concepts & Terminology. Here you use them to make decisions.
Before any move: describe each dataset
Candidates lose storage discussions by naming a database first. Ask these five questions about each dataset instead, and the choice usually falls out:
| Ask | Why it matters |
|---|---|
| How is it read? By exact key, by key plus a range, by arbitrary filters, full text, aggregates, graph hops? | Picks the data model and the indexes |
| How is it written? Peak rate, insert-only or update-in-place, bursty? | Picks the engine, and whether one primary is enough |
| What must stay true across rows? | Needs transactions, and decides which data lives together |
| What does a lost write or a stale read cost? | Picks the replication mode, cache TTLs and backups |
| How big, and for how long? | One node or partitions; retention and cheaper tiers |
Answered for the ride-hailing app, GPS and trips come out as opposites. Current GPS positions are read by location, written 125,000 times a second, have no cross-row rules, cost nothing to lose and are tiny. Trips are read by rider and by driver, written a few hundred times a second, carry money invariants, must never be lost and grow for years.
Size it before you choose
Storage arithmetic is short: bytes per record Γ records per day Γ days kept, then multiply by the copies you keep and add the indexes. Blobs (photos, documents) go to object storage and are sized separately. Here is the app's storage (decimal units, 1 TB = 10ΒΉΒ² bytes; the 50% index overhead is an assumption for this narrow table, not a rule):
| Dataset | Record | Volume | Kept | Raw | Stored | Lives in |
|---|---|---|---|---|---|---|
| Trips | 500 B | 2 million/day | 3 years | 1.1 TB | +50% indexes, 4 copies (primary, 2 standbys, 1 read replica): 6.6 TB | PostgreSQL |
| GPS history | 60 B | 4.32 billion/day (200,000 drivers online on average) | 30 days | 7.8 TB | Object storage replicates for you | Object storage |
| Current positions | 100 B | 500,000 | Latest only | 50 MB | In RAM | In-memory store |
| Driver documents | 2 MB | 6 per driver, 1 million drivers | Account lifetime | 12 TB | Object storage | Object storage |
Two surprises that interviewers like to see you notice: the GPS history is about 260 times bigger per day than the trips (259 GB against 1 GB), and the driver documents outweigh the trips table. Neither belongs in the transactional database.
β οΈ The classic slips: forgetting the copies (three replicas triple the disk), mixing per day and per second (a factor of 86,400), bits and bytes (a factor of 8), and 10Β³ against 2ΒΉβ° (a TiB is about 10% bigger than a TB). State the unit every time you write a number.
Move 1: Place each dataset by its access pattern
Pressure: one database is serving datasets with very different shapes, rates or lifetimes.
"SQL or NoSQL?" is the wrong question
"NoSQL" is not one thing. It is a family of stores that each give up something relational databases offer (joins, ad-hoc queries, multi-row transactions) to get something specific in return. So the useful question is: which of these trade-offs does this dataset want?
Start from the default. A relational database (PostgreSQL, MySQL) gives you any query you can write in SQL, constraints the database enforces, and multi-row transactions. It scales reads with replicas and caches. What it does not do on its own is spread writes across machines; for that you shard it yourself (Move 6) or use a distributed SQL database. As a working assumption, a well-provisioned primary handles a few thousand to low tens of thousands of simple writes a second. Say it as an assumption and offer to benchmark.
| Family | Built for reads like | Gives up | Examples |
|---|---|---|---|
| Relational | Any filter, join or aggregate; constraints and multi-row transactions | Write scaling past one primary needs sharding | PostgreSQL, MySQL |
| Document | Fetch or replace one whole entity with nested, varying fields | Joins are weaker; invariants across documents need multi-document transactions, which cost more | MongoDB, Couchbase |
| Key-value | Get or put by exact key | Any query not by key | Redis, Memcached, DynamoDB |
| Wide-column | All rows under one partition key, sorted by a clustering key; very high write rates | Joins, ad-hoc filters, transactions across partitions; one table per query | Cassandra, ScyllaDB, Bigtable, HBase |
| Search index | Full text, fuzzy matching, facets | It is a derived copy, not the source of truth | Elasticsearch, OpenSearch |
| Columnar warehouse | Aggregates over billions of rows that touch a few columns | Fast single-row reads and updates | BigQuery, Snowflake, ClickHouse |
| Graph | Many-hop traversals: friends of friends, fraud rings | Bulk scans and aggregates | Neo4j, Amazon Neptune |
| In-memory structures | Sub-millisecond counters, sorted sets, geo queries | Capacity is RAM; durability is optional | Redis, Valkey |
| Object storage | Whole blobs by key, cheaply and durably | Queries inside the objects | S3, Google Cloud Storage |
A non-relational store earns its place when one of these is true:
- Every query is by key or partition, you can list them all up front, and the scale is past one primary. You would rather have partitioning built in than shard by hand.
- Writes are huge and simple: append-heavy events, metrics, messages, with no joins.
- The shape fights tables: deeply nested documents with varying fields, graphs traversed many hops deep, full-text search.
- The data is ephemeral: speed matters more than durability and it can be rebuilt.
Designing for a wide-column store: queries first
In a relational schema you model the entities and write queries later. In a wide-column store you do the reverse: one table per query, keyed so that each query reads one partition in order.
Discord has described exactly this. Its messages lived in a single MongoDB replica set, and in late 2015, at about 100 million messages, "the data and the index could no longer fit in RAM and latencies started to become unpredictable" (How Discord Stores Billions of Messages, 2017). In Cassandra the table key became ((channel_id, bucket), message_id):
- Partition key
(channel_id, bucket): one channel's messages from one time window. The bucket covers about 10 days, which keeps each partition comfortably under 100 MB. - Clustering key
message_id: a time-sortable ID, so rows inside a partition are stored in time order. "The latest 50 messages in this channel" usually reads one partition, in order, and stops; a quiet channel walks back through older buckets until it has 50.
The general recipe: the partition key is what you look up by, the clustering key is the order you read in, and a bucket keeps any one partition from growing without limit.
Placing the ride-hailing datasets
driver apps --GPS every 4 s--> location service --GEOADD----> GEO keys by city current positions,
rider apps --cars near me----> location service --GEOSEARCH-> or cell, on a 50 MB, rebuildable
| Redis Cluster
|
+--> log --> object storage GPS history, 30 days
rider apps --book, pay, history--> trip service --> PostgreSQL primary trips, payments:
| the source of truth
+--> 2 standbys take over on failure
+--> read replica trip-history reads
+--> CDC --> warehouse analytics
"Cars near me" becomes a radius query (Redis GEOSEARCH) on data held entirely in RAM. Memory was never the constraint, since 50 MB fits anywhere. The command rate is. A Redis node runs commands on one thread, and a key lives on one node. On a laptop, one node took about 200,000 GEOADDs a second but only about 3,500 radius searches a second when each 2 km search had to scan about 1,600 drivers. The same search over 500 m, which scans about 100 drivers, ran at about 60,000 a second. A single drivers key taking 125,000 updates and 10,000 searches a second would be exactly the hot key Move 6 warns about. So:
- Key the sets by city (
drivers:nyc). With 40 cities, the average key takes about 3,000 updates a second. - Split the busiest cities into geohash cells. If New York has 10% of the drivers, it takes 12,500 updates a second, too many for one key. A search near a cell edge also asks the neighbouring cells.
- Spread the keys over a Redis Cluster, and serve searches from replicas, since positions a few milliseconds stale are fine.
- Search a small radius first and widen it only when too few cars come back.
State the per-node assumption out loud, as you would for the database: one node handles on the order of 100,000 to 200,000 simple commands a second, and far fewer searches. Benchmark yours.
The GPS history, needed only for disputes and analytics, is appended to a log and lands in cheap object storage. The log is the subject of Messaging & Queues. Trips and payments stay in PostgreSQL. The analytics warehouse is fed by change data capture (CDC: streaming the database's committed changes), so year-long reports never compete with trip writes.
Every arrow leaving the source of truth creates a copy that can lag. For each fact, name the store that owns it. Every extra store also adds a failure mode, an on-call runbook and a sync path, so each one must solve a problem the others cannot.
Your turn: place three new datasets.
- In-app chat between rider and driver during a trip, about 6 messages per trip.
- Surge multipliers for a few thousand zones, recomputed every 30 seconds by a pricing job and read on every price quote (3,000 quotes a second at peak).
- Finance wants "revenue per city per hour for the last 2 years" dashboards.
Check your answer
- 2 million trips Γ 6 = 12 million messages a day: about 140 a second on average and about 420 at peak. A PostgreSQL table keyed by
(trip_id, sent_at)handles that. A wide-column store is also defensible if chat grows into a product of its own; the Chat Application (WhatsApp) lesson covers that case. - Tiny, rewritten every 30 seconds, read 3,000 times a second: an in-memory key per zone with a TTL a little longer than 30 seconds, or even a copy inside each app server. The pricing job is the source of truth, so losing the cache costs one recompute. Because the job pushes each new value, readers never recompute it, so this key cannot stampede (Move 5).
- A columnar warehouse fed by CDC. Two-year aggregates on the transactional primary would scan most of a terabyte (730 GB of rows, more with indexes) and compete with trip writes.
Move 2: Index for the query you run most
Pressure: a frequent query filters or sorts, and the plan says "scan".
The trip-history screen runs this on every open:
SELECT id, requested_at, fare_cents
FROM trips
WHERE rider_id = $1
ORDER BY requested_at DESC
LIMIT 20;
With no index, the database reads the whole table to find one rider's rows: 2.19 billion rows and 1.1 TB per request. With an index on (rider_id, requested_at), it walks a B-tree to that rider's newest entry and reads 20 entries.
What a B-tree buys
A B-tree index keeps keys sorted in pages. Each page holds hundreds of keys and points to the pages below it, so each extra level multiplies capacity by that fanout. With about 300 keys per page, four levels cover 300β΄ β 8 billion entries, so a lookup among 2.19 billion rows touches about four pages, and the top levels are almost always cached in memory. Because the keys are sorted, the same structure answers equality (rider_id = 7), ranges (requested_at > ...) and "in order" (ORDER BY ... LIMIT 20), read forward or backward.
Column order is the whole game
An index on (rider_id, requested_at) is sorted by rider first, then by time within each rider. Think of a phone book sorted by last name, then first name:
WHERE rider_id = 7 ORDER BY requested_at DESC LIMIT 20: jump to rider 7, read backward. βWHERE rider_id = 7 AND requested_at >= '2026-09-01': jump to rider 7 at September 1. βWHERE requested_at >= '2026-09-01'alone: trips since September 1 are scattered under every rider, so this index cannot seek. β
The rule is equality columns first, then the one range or sort column. If a query needs a few more columns, a covering index stores them in the index (INCLUDE (fare_cents, id) in PostgreSQL) so the query rarely touches the table. That only works when every selected column is in the index, and PostgreSQL still visits table pages that vacuum has not yet marked all-visible.
The price: every write updates every index
Here is the trade-off measured with SQLite (Python's standard library), on 500,000 trips:
import sqlite3, random, time
def build(n_rows, indexes):
db = sqlite3.connect(":memory:")
db.execute("CREATE TABLE trips (id INTEGER PRIMARY KEY, rider_id INT, "
"driver_id INT, requested_at INT, fare_cents INT)")
for sql in indexes:
db.execute(sql)
rnd = random.Random(7)
rows = [(i, rnd.randrange(50_000), rnd.randrange(10_000), i, rnd.randrange(500, 9_000))
for i in range(n_rows)]
start = time.perf_counter()
db.executemany("INSERT INTO trips VALUES (?, ?, ?, ?, ?)", rows)
db.commit()
return db, time.perf_counter() - start
QUERY = ("SELECT id, requested_at, fare_cents FROM trips "
"WHERE rider_id = ? ORDER BY requested_at DESC LIMIT 20")
def plan(db):
return [row[-1] for row in db.execute("EXPLAIN QUERY PLAN " + QUERY, (42,))]
def time_query(db, repeats=200):
start = time.perf_counter()
for r in range(repeats):
db.execute(QUERY, (r,)).fetchall()
return (time.perf_counter() - start) / repeats
N = 500_000
plain, t_plain = build(N, [])
indexed, t_indexed = build(N, [
"CREATE INDEX trips_rider_time ON trips (rider_id, requested_at)",
"CREATE INDEX trips_driver_time ON trips (driver_id, requested_at)",
"CREATE INDEX trips_time ON trips (requested_at)",
])
print("plan without index:", plan(plain))
print("plan with index: ", plan(indexed))
print(f"query: {time_query(plain)*1e3:.2f} ms -> {time_query(indexed)*1e3:.3f} ms")
print(f"insert {N:,} rows: {t_plain:.2f} s with 0 extra indexes, {t_indexed:.2f} s with 3")
One run on a laptop printed:
plan without index: ['SCAN trips', 'USE TEMP B-TREE FOR ORDER BY']
plan with index: ['SEARCH trips USING INDEX trips_rider_time (rider_id=?)']
query: 13.38 ms -> 0.015 ms
insert 500,000 rows: 0.28 s with 0 extra indexes, 1.04 s with 3
Your timings will differ; the ratios are the point. The query got about 900 times faster, and the inserts got about 3.7 times slower, because every insert now also finds its place in three more sorted structures. Updates pay too. In MySQL an update touches only the indexes whose columns change. In PostgreSQL an update that changes any indexed column, or whose new row version does not fit on the same page, adds an entry to every index on the table. Indexes also take disk and memory, often a large fraction of the table itself.
So index for the queries you run often and name the write cost when you propose one. Three more habits:
- An index on a column with a handful of values (a
statuswith 5 states) rarely helps on its own. Use a partial index such asWHERE status = 'searching', or put the column inside a composite index. - On a big live table, build new indexes without blocking writes (
CREATE INDEX CONCURRENTLYin PostgreSQL). See the failure drills below. - Read the plan (
EXPLAIN) instead of guessing.
Your turn: the driver app shows "my earnings this week": SELECT sum(fare_cents) FROM trips WHERE driver_id = ? AND finished_at >= ?. Which index, and why is (finished_at, driver_id) worse?
Check your answer
(driver_id, finished_at) INCLUDE (fare_cents). It jumps to one driver's week and sums fares, usually without touching the table. With (finished_at, driver_id) the range column comes first, so the index can only seek to "this week" and must then read every driver's trips for the week and filter them: roughly as many times more entries as there are active drivers.
B-tree or LSM-tree: where the write cost goes
A B-tree updates pages in place, so writes land at random places on disk while reads stay cheap and predictable. That suits most mixed workloads. Write-heavy stores such as Cassandra, ScyllaDB and RocksDB use a log-structured merge tree (LSM-tree) instead:
- A write is appended to a commit log and inserted into an in-memory sorted table (the memtable).
- A full memtable is flushed to disk as an immutable sorted file (an SSTable).
- Background compaction merges SSTables, dropping overwritten and deleted values.
Writes become cheap sequential appends. The costs move elsewhere: compaction rewrites data repeatedly (write amplification), a read may check the memtable and several SSTables (read amplification; Bloom filters let it skip most files), and a delete is a tombstone that occupies space until compaction removes it.
Say it in the interview: "This workload is append-heavy with no joins, so an LSM-based store fits. I'm accepting compaction overhead and slower reads of old data." Then move on. The engine rarely decides an interview by itself.
Move 3: Keep invariants true under concurrency
Pressure: several writes must succeed together, or two requests race for the same row.
The double-dispatch bug
Two matching workers pick driver 42 for two different riders in the same millisecond. Both read status = 'available', and both write status = 'on_trip'. Now one driver has two riders. Nothing crashed, and every statement was individually correct. The bug lives in the gap between the read and the write.
Close the gap by making the check and the write one statement:
UPDATE drivers
SET status = 'on_trip', trip_id = 9001
WHERE id = 42 AND status = 'available';
-- 1 row updated: you got the driver. 0 rows: someone else did, so try the next one.
The second worker's UPDATE waits for the first one's row lock. Under PostgreSQL's default Read Committed level, once the first commits, PostgreSQL re-evaluates the WHERE clause against the new row version, finds status = 'on_trip' and updates 0 rows. Under Repeatable Read the loser gets a serialization error and retries instead. Other stores offer the same move: a DynamoDB condition expression, a Cassandra lightweight transaction (... IF status = 'available', built on Paxos, with several extra round trips), a Redis Lua script.
The lost update, reproduced
The same gap appears in any read-modify-write done in application code. Here two app servers charge one rider's wallet at the same time:
import sqlite3
URI = "file:wallets?mode=memory&cache=shared"
setup = sqlite3.connect(URI, uri=True, isolation_level=None) # keeps the shared DB alive
setup.execute("CREATE TABLE wallets (rider_id INT PRIMARY KEY, balance_cents INT)")
setup.execute("INSERT INTO wallets VALUES (7, 10000)") # 100.00
a = sqlite3.connect(URI, uri=True, isolation_level=None) # app server A
b = sqlite3.connect(URI, uri=True, isolation_level=None) # app server B
def balance(conn):
return conn.execute("SELECT balance_cents FROM wallets WHERE rider_id = 7").fetchone()[0]
# Read-modify-write in application code: charge 30.00 on A and 50.00 on B at the same time
seen_a = balance(a) # A reads 10000
seen_b = balance(b) # B reads 10000 too
a.execute("UPDATE wallets SET balance_cents = ? WHERE rider_id = 7", (seen_a - 3000,))
b.execute("UPDATE wallets SET balance_cents = ? WHERE rider_id = 7", (seen_b - 5000,))
print("read-modify-write:", balance(a)) # 5000: A's charge vanished
# The same two charges, each as one conditional statement
setup.execute("UPDATE wallets SET balance_cents = 10000 WHERE rider_id = 7")
for conn, charge in ((a, 3000), (b, 5000)):
cur = conn.execute("UPDATE wallets SET balance_cents = balance_cents - ? "
"WHERE rider_id = 7 AND balance_cents >= ?", (charge, charge))
print("charged" if cur.rowcount == 1 else "declined", charge)
print("atomic update: ", balance(a)) # 2000
It prints 5000, then two "charged" lines, then 2000. B computed its new balance from a value that was already out of date, so A's 30.00 charge disappeared. Letting the database do the arithmetic removes the gap, and the balance_cents >= ? condition also enforces "never negative". The other standard fixes are a pessimistic lock (SELECT ... FOR UPDATE), optimistic concurrency (a version column and WHERE version = 17, retrying on 0 rows), or putting the read and the write in one transaction at Repeatable Read or Serializable with a retry loop. In autocommit mode, as in the demo, each statement is its own transaction, and no isolation level can see the gap.
ACID in one breath each
- Atomicity: all of a transaction's writes commit, or none do. Finishing a trip inserts the trip result and two ledger rows (rider debit, driver credit) in one transaction, so a crash can never leave money half-moved.
- Consistency: the database enforces the rules you declared (constraints, foreign keys, unique indexes), and your transactions supply the rest. This is not the C in CAP; see Key Concepts & Terminology.
- Isolation: how much concurrent transactions see of each other. It is a dial, not a switch.
- Durability: once the commit is acknowledged, it survives a crash.
Durability comes from the write-ahead log (WAL). The order matters:
COMMIT
|
v
append the commit record to the write-ahead log (sequential write)
|
v
flush the log to disk (fsync) --> reply "committed" to the client
later, in the background: write the changed data pages (checkpoint)
after a crash: replay the log from the last checkpoint
The client waits for the log flush, not for the data pages, which is why commits are fast. You can trade some durability for latency: PostgreSQL's synchronous_commit = off, set per transaction, acknowledges before the flush, so a crash can lose the last fraction of a second of commits (without corrupting anything). Durability on one disk still dies with that machine; Move 4 extends it to another one.
Isolation levels at interview depth
| Anomaly | What it looks like | Read Committed | Snapshot (PostgreSQL Repeatable Read) | Serializable |
|---|---|---|---|---|
| Dirty read | You see another transaction's uncommitted write | Prevented | Prevented | Prevented |
| Non-repeatable read | The same row read twice gives two values | Possible | Prevented | Prevented |
| Phantom | A repeated range query returns new rows | Possible | Prevented in PostgreSQL | Prevented |
| Lost update | The wallet charge, read and write in one transaction | Possible | The second writer gets an error and retries | Prevented |
| Write skew | Two transactions check a rule, then each writes a different row | Possible | Possible | Prevented |
Defaults differ, so name your engine. PostgreSQL defaults to Read Committed; MySQL's InnoDB defaults to Repeatable Read. InnoDB's plain SELECT reads a snapshot, but an UPDATE reads the latest row, so an application-side read-modify-write can still lose updates there. Use atomic updates or SELECT ... FOR UPDATE. Serializable is the safest level, but it costs throughput and forces you to retry aborted transactions.
Write skew deserves an example because no row lock catches it. Rule: a driver has at most one active trip. Two transactions each run SELECT count(*) FROM trips WHERE driver_id = 42 AND status IN ('matched', 'in_progress'), both see 0, and both insert a trip. They wrote different rows, so nothing conflicted. Fixes: Serializable (PostgreSQL detects the pattern and aborts one), lock the driver's row first, or, best, let a constraint hold the rule under any isolation level:
CREATE UNIQUE INDEX one_active_trip_per_driver
ON trips (driver_id)
WHERE status IN ('matched', 'in_progress');
-- the second concurrent INSERT now fails with a unique-violation error
Where transactions stop
An ordinary transaction covers one database. There is no transaction spanning Redis and PostgreSQL. Two-phase commit can span two relational databases, but it blocks when its coordinator fails, so services avoid it (see Advanced Topics & Final Prep). Choosing a "CP" database does not change that: CAP is about replicas during a network partition, not about making two services' writes atomic. So keep each invariant that must never break inside one database, ideally in one row or one shard, and connect services with the outbox and saga patterns from Advanced Topics & Final Prep.
Say it in the interview: "The invariants that must never break are: no driver on two trips, no negative wallet, ledger rows sum to zero. Each one is enforced inside PostgreSQL: by a conditional write, by a constraint, or for the ledger by writing both rows in one transaction. Everything else can be eventually consistent."
Your turn: a promo code may be redeemed 1,000 times in total. The code does SELECT count(*) FROM redemptions WHERE code = ? and inserts a redemption when the count is under 1,000, at Read Committed. What goes wrong at launch, and what is the fix?
Check your answer
Many requests read 999 at the same moment and all insert, so the code is redeemed well over 1,000 times. This is write skew: each insert is a different row. Fix: keep a counter row and claim with one conditional statement, UPDATE promos SET remaining = remaining - 1 WHERE code = ? AND remaining > 0, and insert the redemption in the same transaction only if one row was updated. At national-launch rates that one row becomes a hot key, and Move 6 shows how to split it.
Move 4: Replicate to survive failures and to spread reads
Pressure: a machine will die and the business cannot wait for a restore; or reads outgrow one machine.
writes reads that can be a little stale
| |
v v
+-----------+ async log stream +--------------+
| primary |-------------------------->| read replica |
+-----------+ +--------------+
|
| sync log stream: any 1 of the 2 standbys must confirm each commit
+------+-------+
v v
+-----------+ +-----------+
| standby A | | standby B | one is promoted if the primary dies
+-----------+ +-----------+
The primary streams its log to the other copies. The only real question is when the primary says "committed":
- Asynchronous: after its own log flush. Commits are fast. But if the primary dies, anything not yet shipped is gone, even though clients were told it committed. At 200 ms of lag and 350 trip writes a second, about 70 acknowledged writes vanish in a failover.
- Synchronous: after a standby also confirms. A single failure loses nothing, but every commit pays a network round trip: about a millisecond or two inside a region, tens of milliseconds across regions. If the only synchronous standby is down, commits wait. So use a quorum, for example PostgreSQL's
synchronous_standby_names = 'ANY 1 (s1, s2)'or MySQL's semi-synchronous replication, so either standby can confirm. (By default, MySQL semi-sync silently falls back to asynchronous replication after a timeout, so monitor it.)
For the ride-hailing app: payments are money, so run two standbys in other availability zones with ANY 1, so either one confirms each commit and neither can block writes alone, plus one asynchronous read replica for trip history. That is the four copies the storage table counted. A few milliseconds per commit is cheap at 350 writes a second.
Before a failover, also make sure the old primary cannot keep accepting writes if it was only partitioned rather than dead. That is fencing, covered with leader election in Advanced Topics & Final Prep.
Replication lag: the bugs users notice
Lag is usually milliseconds. Under heavy writes, long-running queries on a replica or network trouble it can grow to seconds or minutes. Three symptoms:
- Missing your own write. A rider adds a card, the next page reads a replica 300 ms behind, the card is not there, and the rider adds it again. Fixes: read a user's own recent changes from the primary (for example for a minute after they write); or remember the log position of the write (a PostgreSQL LSN or MySQL GTID) and read from a replica only once it has replayed that far.
- Time going backward. Two refreshes hit two replicas with different lag, so the trip shows "finished" and then "in progress". Fix: pin each user's reads to one replica, for example by hashing the user ID.
- Other people's data a bit old. A driver's rating that is 2 seconds old is usually fine. Say so explicitly instead of paying for consistency nobody needs.
Monitor lag, and take a replica out of the read pool when it falls further behind than your staleness budget.
Leaderless replication: quorums instead of a primary
Cassandra and ScyllaDB follow the design of Amazon's Dynamo paper. There is no primary: any node can coordinate a request, a write is sent to all N replicas and succeeds when W of them acknowledge, and a read asks R of them. When W + R > N, every read quorum overlaps every write quorum, so a read reaches at least one replica with the latest acknowledged write. With N = 3, QUORUM means 2.
What you give up:
- Conflicts are settled by timestamp. Cassandra uses last-write-wins: of two concurrent writes to the same cell, one is silently discarded, and a skewed clock can make the "wrong" one win. For compare-and-set you need lightweight transactions and their extra round trips.
- Failures force a choice per query. Take 6 nodes with replication factor 3, and let two neighbouring nodes die. Some key ranges now have one live replica:
QUORUMreads and writes for those keys fail, whileONEstill succeeds and may return stale data. You choose availability or consistency per request, which is CAP in practice. Even with no failures,QUORUMwaits for more replicas thanONE, so you trade latency for consistency (the "else" half of PACELC).
Multi-leader replication, with a writable primary in each region, has the same conflict problem across regions. Use it only when users in several regions must write locally and you have a conflict rule you can defend.
Your turn: the primary is in us-east and an asynchronous replica in eu-west serves European reads. Lag is normally 150 ms and occasionally 20 seconds. A rider in Paris finishes a trip and immediately opens their history. What might they see, and what is the cheapest fix that does not send every European read across the Atlantic?
Check your answer
The history can be missing the trip that just ended, for up to 20 seconds on a bad day. The cheapest fix uses what the client already has: the "trip finished" response contains the trip, so the app shows it at the top of the history until the replica catches up. Alternatives: route that rider's history reads to the primary for a minute after a trip ends (a few cross-Atlantic reads per trip, not all of them), or read from the replica only once it has replayed the trip's log position. Synchronous replication to Europe would also fix it, but it would add a transatlantic round trip to every commit.
Move 5: Cache what is read far more often than it changes
Pressure: the same data is read over and over, and those reads dominate the database.
During a trip the rider's app polls the trip card: driver name, photo, car, plate, rating. Peak trips start at about 69 a second, and each lasts about 25 minutes from match to drop-off, so by Little's law about 100,000 trips are in progress. With one poll every 5 seconds, that is 20,000 reads a second of data that rarely changes: the rating moves after rated trips, the rest a few times a year. The complete set of driver cards (1 million drivers Γ 2 KB = 2 GB) fits in one cache node.
Hit ratio is a load multiplier
| Cache hit ratio | Reads reaching the database |
|---|---|
| 99% | 200/s |
| 95% | 1,000/s |
| 90% | 2,000/s |
| 80% | 4,000/s |
| 0% (cold cache) | 20,000/s |
Dropping from 95% to 90% sounds small and doubles the database load. A cold cache after a restart sends all 20,000 reads a second at a database sized for 1,000. Size the database for the miss rate you can actually get into, and plan the warm-up.
Read and write patterns
| Pattern | Read path | Write path | What you give up |
|---|---|---|---|
| Cache-aside (lazy loading) | App checks the cache; on a miss it reads the database and fills the cache | App writes the database, then deletes the cache key | A miss after every change; a race can leave a stale value until the TTL |
| Read-through | The cache library loads from the database on a miss | Combined with any write pattern | Same as cache-aside, with the logic inside the cache layer |
| Write-through | Reads find a warm cache | The write goes to the cache and the database before it is acknowledged | Slower writes; caches data nobody reads |
| Write-behind (write-back) | Reads hit the cache | Acknowledge after the cache write; flush to the database later, often batched | Writes since the last flush are lost if the cache node dies |
Cache-aside is the default answer in interviews and in practice. Here are its read path, its write path, and the fix for a stampede:
import threading, time
drivers = {42: {"name": "Ana", "car": "Blue Corolla"}} # the database: source of truth
db_queries = 0
count_lock = threading.Lock()
def db_read(driver_id):
global db_queries
with count_lock:
db_queries += 1
time.sleep(0.05) # a 50 ms query
return dict(drivers[driver_id])
cache = {} # key -> (value, expires_at)
TTL = 300 # the staleness bound, in seconds
def cached(key):
hit = cache.get(key)
return hit[0] if hit and hit[1] > time.monotonic() else None
def get_driver_naive(driver_id): # cache-aside read path
key = f"driver:{driver_id}"
value = cached(key)
if value is None:
value = db_read(driver_id) # every concurrent miss hits the DB
cache[key] = (value, time.monotonic() + TTL)
return value
locks, locks_guard = {}, threading.Lock()
def get_driver(driver_id): # same read path, with single-flight
key = f"driver:{driver_id}"
value = cached(key)
if value is not None:
return value
with locks_guard:
lock = locks.setdefault(key, threading.Lock())
with lock: # one loader per key; the rest wait
value = cached(key) # re-check: the loader may have filled it
if value is None:
value = db_read(driver_id)
cache[key] = (value, time.monotonic() + TTL)
return value
def update_car(driver_id, car): # cache-aside write path
drivers[driver_id]["car"] = car # 1. write the source of truth
cache.pop(f"driver:{driver_id}", None) # 2. then delete the cached copy
def burst(getter, requests=200):
global db_queries
cache.clear()
db_queries = 0
threads = [threading.Thread(target=getter, args=(42,)) for _ in range(requests)]
for t in threads:
t.start()
for t in threads:
t.join()
return db_queries
print("200 concurrent misses, naive: ", burst(get_driver_naive), "DB queries")
print("200 concurrent misses, single-flight:", burst(get_driver), "DB query")
update_car(42, "Red Prius")
print("next read after the update:", get_driver(42)["car"])
200 concurrent misses, naive: 200 DB queries
200 concurrent misses, single-flight: 1 DB query
next read after the update: Red Prius
Why delete instead of update on write
Suppose the write path set the new value in the cache instead. Two writers change the same driver's car, first to Green and then to Red, and the database ends at Red. Their two cache sets travel separately, so they can arrive in the opposite order: Red, then Green. With no TTL, the cache now says Green forever. Two deletes, in any order, leave the same state: empty, and the next read loads Red. Facebook's paper Scaling Memcache at Facebook (NSDI 2013) puts it in one line: "We choose to delete cached data instead of updating it because deletes are idempotent."
The race that survives delete-on-write
Deleting is not a complete fix. Trace a slow reader and a writer:
| Step | Reader (cache miss) | Writer | Database | Cache |
|---|---|---|---|---|
| 1 | Cache miss | Blue | empty | |
| 2 | Reads the database: Blue | Blue | empty | |
| 3 | Updates the database to Red | Red | empty | |
| 4 | Deletes the key (nothing there) | Red | empty | |
| 5 | Sets the cache to Blue, TTL 300 s | Red | Blue |
The cache now serves Blue for up to 300 seconds. It takes a slow reader and a write landing inside its window, which is rare per key but routine across millions of keys. Always set a TTL: it is the bound on how long any bug or race can serve stale data. For tighter guarantees, Facebook's paper describes leases. A miss gets a token, and a delete invalidates any outstanding token, so the late set in step 5 is rejected.
Eviction and sizing
A cache is smaller than the data, so something must be evicted:
- LRU (least recently used) keeps what was touched recently. It is the usual default.
- LFU (least frequently used) keeps what is popular over time, so a one-off scan cannot flush the popular keys.
- TTL expiry removes data by age, whatever the memory pressure.
β οΈ Redis's defaults do not make a cache. maxmemory is 0 (no limit) on 64-bit builds, so Redis grows until the operating system kills it, and the default maxmemory-policy is noeviction, so once a limit is reached writes return errors instead of evicting anything. For a pure cache, set maxmemory and allkeys-lru or allkeys-lfu. Memcached evicts by LRU (per slab class) by default.
Size a cache by its working set: how many keys receive most of the reads. Then measure the hit ratio at that size.
Stampedes: when a hot key expires
The driver app's demand heat map for New York is cached cache-aside. About 50,000 drivers online there poll it every 20 seconds, so it is read 2,500 times a second, and rebuilding it aggregates recent ride requests for 2 seconds of database time. When it expires, every request during those 2 seconds misses: 5,000 rebuilds of the same value, all hitting the database at once. This is a cache stampede (also called a thundering herd or dogpile). Fixes:
- Single-flight (request coalescing): one request per key recomputes and the rest wait for its result. The demo above goes from 200 database queries to 1. It works per process, so 40 app servers still send up to 40 queries. To get down to one, take a short lock or lease in the cache itself. Facebook reports that its leases cut the peak database query rate for keys prone to this from 17,000 to 1,300 a second. Discord has described data services that coalesce identical concurrent reads and route requests by consistent hashing, so identical requests arrive at the same service instance.
- Serve stale while refreshing: keep the old value past its soft expiry and let one request refresh it in the background. Users see data a few seconds old instead of waiting.
- Refresh early: recompute popular keys before they expire, with some randomness so that not every server does it at the same moment.
- Jitter the TTLs: add a random Β±10% so that keys cached together (after a deploy or a warm-up) do not all expire together. β οΈ Jitter does nothing for a single hot key.
When the cache is the database
Know what your in-memory store keeps after a crash. Redis's own defaults (see Redis persistence):
- Snapshots (RDB) are on by default: after 3,600 s if at least 1 key changed, 300 s if 100 changed, 60 s if 10,000 changed. A crash loses everything written since the last snapshot, which can be minutes (up to an hour when few keys change). A clean shutdown writes a final snapshot.
- The append-only file (AOF) is off by default. With it on and the default
appendfsync everysec, a power loss costs about one second of writes. - Many cache deployments turn persistence off entirely, and then a restart loses everything. Redis replication is asynchronous, so a failover can also lose recent writes.
Decide by asking what loss costs:
- Driver positions: nothing. Every driver reports again within 4 seconds, so no persistence is needed. There is one trap, though. Each city's (or cell's) GEO set is a single sorted-set key, and a TTL expires the whole key, not one driver. To drop drivers who went silent, keep a second sorted set per key, scored by last-seen time, and periodically read the members older than 30 seconds (
ZRANGEwithBYSCORE) andZREMthem from both sets, or write into per-minute keys that expire as a whole. - Sessions: losing them logs everyone out. Some products accept that; say which kind yours is.
- Carts, orders, balances: they must survive. The database is the source of truth and the cache is only a cache. Running Redis as the primary store for them needs AOF, replicas and a loss window you state out loud.
Say it in the interview: "Redis is the cache and Postgres is the record. If Redis disappears we get slower, not wrong. The follow-up is the cold start, so I'd warm the hot keys and cap concurrent misses."
Move 6: Partition when one primary is not enough
Pressure: the data or the write rate has outgrown one primary, even after the moves above.
Sharding comes last because it is the hardest move to undo. Every query now needs the shard key, joins and transactions across shards go away, and you operate N databases. Climb the ladder first: fix queries and indexes, scale the machine up, cache hot reads, add replicas. Splitting a table into partitions on one server (PostgreSQL declarative partitioning by month, say) is also useful. Dropping a month becomes instant instead of deleting tens of millions of rows, but it adds no write capacity.
Now grow the ride-hailing app 10Γ: 20 million trips a day, about 3,500 writes a second at peak, and about 11 TB of trips over three years, about 16 TB with indexes. A big primary might still take the writes. The size is what hurts: at 500 MB/s, restoring 16 TB from a backup takes about 9 hours, and index builds and major-version upgrades slow down the same way. That is a real reason to split.
Range, hash or directory
| Scheme | How a key finds its shard | Good at | Bad at |
|---|---|---|---|
| Range | Each shard owns a contiguous range of keys | Scans across keys ("all trips from 9:00 to 10:00") | Increasing keys (timestamps, auto-increment IDs) send every new write to the last shard |
| Hash | A stable hash of the key picks the shard | Spreading many keys evenly | Range scans across keys become queries to every shard |
| Directory | A lookup table maps key (or tenant) to shard | Moving one big tenant by itself | The directory is a critical dependency; cache it |
Two corrections to common folklore. Sequential IDs are only a hot spot under range partitioning; hashed, they spread evenly. And hashing spreads keys, not load: it cannot split a single hot key (see below). You can also combine schemes. Hash the partition key to spread riders, and sort by a clustering key inside each partition, which is the Discord pattern from Move 1.
Do not route with hash % N
import zlib
from collections import Counter
def stable_hash(key: str) -> int:
return zlib.crc32(key.encode()) # same result in every process, unlike hash()
keys = [f"rider:{i}" for i in range(100_000)]
# Scheme 1: server = hash % number_of_servers
moved = sum(stable_hash(k) % 4 != stable_hash(k) % 5 for k in keys)
print(f"hash % servers, 4 -> 5 servers: {moved / len(keys):.0%} of keys move")
# Scheme 2: 1,024 fixed logical partitions, plus a small partition -> server table
PARTITIONS = 1024
def partition_of(key):
return stable_hash(key) % PARTITIONS # a key's partition never changes
four = {p: p % 4 for p in range(PARTITIONS)}
five = {p: (4 if p % 5 == 4 else s) for p, s in four.items()} # new server takes every 5th
moved = sum(four[partition_of(k)] != five[partition_of(k)] for k in keys)
print(f"logical partitions, 4 -> 5: {moved / len(keys):.0%} of keys move")
print("partitions per server after:", sorted(Counter(five.values()).items()))
hash % servers, 4 -> 5 servers: 80% of keys move
logical partitions, 4 -> 5: 20% of keys move
partitions per server after: [(0, 205), (1, 205), (2, 205), (3, 205), (4, 204)]
With hash % N, adding a fifth server moves 80% of the data, and a key stays put only when both remainders happen to agree. With many fixed logical partitions and a small table mapping partitions to servers, the new server takes an even 20% share and nothing else moves. Consistent hashing achieves the same without the table; it is covered in Advanced Topics & Final Prep. Two habits: create far more logical partitions than servers on day one, and never route with a language's built-in hash() (Python randomizes string hashes per process).
Choosing the shard key
A good shard key meets four tests:
- The most frequent queries carry it, so they touch one shard.
- The invariants that need transactions live under one key value.
- It has many distinct values, and no single value carries a large share of the load.
- It does not only increase, if you partition by range.
Shard trips by rider_id and rider history hits one shard, but "driver earnings" must ask every shard. Fix: keep a second copy of each trip keyed by driver_id, updated asynchronously, and accept that it lags. Sharding by city is the tempting mistake: there are few values, the biggest city carries a large share of trips, and its shard runs hot while the others idle.
Slack has described the tenant version of this (Scaling Datastores at Slack with Vitess, 2020). Each MySQL shard held all the data for thousands of workspaces. As individual customers grew, "their designated shard reached the largest available hardware", and Slack could not spread those customers' load across the fleet, which left "a few hot spots in our database tier". Vitess let them shard message data by channel ID instead, which spreads the load more evenly. Sharding by tenant is great for isolation and transactions until one tenant outgrows a shard.
Hot keys: when hashing cannot help
A hot key is a single key with more traffic than one shard can serve. Every request for it lands on the same shard, so adding shards does not help. First ask: is it hot for reads or for writes?
- Read-hot (a celebrity's profile, a viral product, New York's demand heat map): cache it with single-flight, keep a copy in each app server's memory for a second or two, add replicas, and put public content on a CDN. Do not salt it. Splitting a read-hot key into sub-keys makes every read gather all the pieces, which is more load, not less. Feeds have their own version of this problem, covered in Twitter/Instagram News Feed.
- Write-hot (a global counter, the stock row of a flash-sale item, a promo code's remaining count): split the key into K sub-keys that live on different partitions (write sharding) and combine them on read. Or aggregate in memory and write once a second (write-behind: you accept losing up to a second). Or queue the writes and apply them serially.
- Both (Move 1's live driver positions): split along the access pattern. Every search is local, so keys by city and geohash cell spread both the updates and the searches.
DynamoDB documents a hard number for this: each partition delivers at most 3,000 read units and 1,000 write units a second, and one write unit is one write of up to 1 KB. So a single item takes at most about 1,000 small writes a second. A like counter receiving 5,000 increments a second needs at least 5 sub-keys, and 10 for headroom. A read then sums 10 items, or reads a total cached for a second.
Discord hit read-side hot partitions in production. By its 2023 account, "one channel and bucket pair received a large amount of traffic" and the node's latency climbed. Coalescing identical reads in its data services was part of the fix, alongside moving from 177 Cassandra nodes to 72 ScyllaDB nodes.
Queries and transactions across shards
- Scatter-gather: send the query to every shard and merge the answers. The latency is the slowest shard's, and the load is multiplied by the shard count. It is fine for rare admin queries and bad on a hot path.
- Secondary indexes: a local index lives on each shard, so a query by the indexed column must ask every shard. A global index is partitioned by the indexed value, so each write also updates another partition, usually asynchronously.
- Transactions across shards: design them away with the shard key. When you cannot, use a saga, from Advanced Topics & Final Prep.
Your turn: trips are sharded by hash of rider_id across 16 shards. A new dispatcher job asks every minute for "rides scheduled for pickup in Chicago between 7:00 and 7:15 tomorrow". Do you need a second copy of the trips, keyed by city and time?
Check your answer
Probably not; an index is enough. Scatter-gather across 16 shards once a minute is 16 small indexed queries a minute (with an index on (city_id, pickup_at) on each shard), which is trivial. Build a separate copy keyed by (city, pickup slot) only if this query moves onto a hot path, for example if riders' apps start polling it. Knowing when not to add machinery is part of the answer.
Failure drills: what breaks and what you say
Interviewers probe storage with "and then what?" Have the answer for each event, including what it costs:
| Event | What happens | Mitigation | What you give up |
|---|---|---|---|
| Primary dies, async replicas only | A replica is promoted; writes in the lag window are lost (about 70 at 200 ms lag) | Synchronous or quorum standby for money data; fence the old primary | Commit latency |
| A replica lags 30 s | Users miss their own writes; pages go backward | Read-your-writes routing; pin sessions; eject lagging replicas | More load on the primary |
| Cache cluster restarts empty | Hit ratio 0%: 20,000 reads/s instead of 1,000 | Warm hot keys, single-flight, cap concurrent misses | Slower recovery |
| One key goes viral | One shard saturates while the others idle | Cache and replicate (read-hot) or split the key (write-hot) | Complexity; slower reads of split keys |
Someone runs DELETE without a WHERE |
Every replica deletes too | Point-in-time restore from backups | Restore time |
| A whole region fails | Everything in that region is gone | Cross-region replica or backups | Latency or cost, plus a lag-sized RPO |
Replicas are not backups: RPO and RTO
RPO (recovery point objective) is how much data, measured in time, you can afford to lose. RTO (recovery time objective) is how long you can be down. Different tools cover different failures:
| Tool | Protects against | RPO | RTO |
|---|---|---|---|
| Synchronous standby | Machine or zone failure | Zero | Seconds to a minute or two |
| Asynchronous replica | Machine failure | The lag at the moment of failure | Seconds to a minute or two |
| Point-in-time recovery (base backup + archived log) | Bad deploys, accidental deletes, corruption | How often the log is archived | Minutes to hours, growing with size |
| Daily snapshot only | Losing everything | Up to 24 hours | Hours |
A replica faithfully replicates your mistakes, so a DELETE or a corrupting bug reaches every copy within milliseconds. Only a backup lets you go back in time. On AWS RDS, point-in-time restore creates a new DB instance, and RDS uploads transaction logs to S3 every five minutes. So the newest restorable point can be about five minutes old, and the restore takes minutes to hours. You need both: a standby for machine failure and backups for human and software failure. A restore you have never tested is only a hope.
Changing a big table without an outage
"Add a column to a 500-million-row table in production" is a classic senior probe. The modern answer is more precise than "it locks the table for hours":
- Adding the column is usually instant. In PostgreSQL, adding a nullable column, or (since PostgreSQL 11) a column with a constant default, only changes metadata. MySQL 8.0 has
ALGORITHM=INSTANTfor the same case. - The danger is the lock queue.
ALTER TABLEneeds a brief exclusive lock. If a long-running transaction holds any lock on the table (even from a plainSELECT), theALTERwaits, and every new query, even a simpleSELECT, queues behind the waitingALTER. A two-second change becomes an outage. Set alock_timeoutand retry. SET NOT NULLscans the whole table while holding that exclusive lock. Add the rule as a constraint that is not yet validated, then validate it under a weaker lock that lets reads and writes continue.
SET lock_timeout = '2s'; -- give up and retry rather than queue all traffic behind us
-- 1. expand: metadata-only change on PostgreSQL 11+
ALTER TABLE riders ADD COLUMN home_city_id bigint;
-- 2. deploy application code that sets home_city_id on every INSERT and UPDATE;
-- from now on no new row is NULL, so the backfill below can finish
-- 3. backfill old rows in small batches from a job, until no NULLs remain
UPDATE riders SET home_city_id = signup_city_id
WHERE id BETWEEN 1 AND 10000 AND home_city_id IS NULL;
-- 4. enforce without a long exclusive lock
ALTER TABLE riders ADD CONSTRAINT home_city_present
CHECK (home_city_id IS NOT NULL) NOT VALID; -- instant; checks new writes only
ALTER TABLE riders VALIDATE CONSTRAINT home_city_present; -- scans; reads and writes continue
ALTER TABLE riders ALTER COLUMN home_city_id SET NOT NULL; -- no scan: the CHECK proves it
Order matters. A NOT VALID constraint is still checked on every new insert and update, so add it only after the application writes the column. Added earlier, it starts rejecting signups.
(PostgreSQL 18 can also add a NOT NULL constraint as NOT VALID directly.) Build new indexes with CREATE INDEX CONCURRENTLY, which does not block writes but is slower and cannot run inside a transaction. On MySQL, tools such as gh-ost and pt-online-schema-change copy rows into a shadow table and swap it in.
For renames and type changes, use expand and contract: add the new column, write to both, backfill, switch reads, stop writing the old column, and finally drop it. Each step is deployable on its own and can be rolled back.
Final round: no label on the problem
Real prompts do not say which move they want. Describe each dataset, find the pressure, then pick the move and name its cost.
Challenge 1: the sneaker drop
A limited sneaker release: 5,000 pairs, 400,000 users pressing Buy in the first 10 seconds after 10:00:00. There must be no overselling, at most one pair per user, and an answer within one second.
Before peeking:
- Where is the hot key, and is it read-hot or write-hot?
- How do you guarantee no overselling?
- How do you enforce one pair per user?
- What do the 395,000 users who miss out touch?
Hint
400,000 attempts in 10 seconds is 40,000 a second on one stock row. If each commit holds the row lock for about 1 ms, one row serializes to about 1,000 purchases a second.
Check your answer
- The stock count is write-hot. Split it: 50 sub-rows of 100 pairs each, with each buyer hashed to one and moved on to another when it is empty. Now 50 rows take conditional decrements in parallel. The product page is read-hot as well; cache it with single-flight. The stock shown there can be approximate.
- No overselling: each claim is
UPDATE stock SET remaining = remaining - 1 WHERE drop_id = ? AND bucket = ? AND remaining > 0and succeeds only if one row changed. An alternative is an in-memory atomic decrement as a gate in front of the database (a single-threaded counter handles tens of thousands of decrements a second), followed by a durable order write, plus a reconciliation job that returns units when the order write fails. - One per user: a unique constraint on
(drop_id, user_id)in the orders table. It holds under any isolation level. - Sold-out requests should hit only the gate or a cached "sold out" flag, not the database. Name the cost: sub-rows make "exactly how many are left?" a sum over 50 rows, and a buyer may be told "sold out" a moment before a cancelled pair returns to a bucket.
Challenge 2: a game leaderboard
A mobile game has 5 million players. At peak there are 20,000 score updates a second, 30,000 "top 100" page loads a second and 20,000 "what's my rank?" lookups a second. Scores come from match results that are already stored durably. The season lasts 30 days.
Check your answer
Store: an in-memory sorted set. ZADD updates a score in O(log n), ZRANGE board 0 99 REV returns the top 100, and ZREVRANK returns a player's rank in O(log n). Assuming 100 to 200 bytes per member including overhead, 5 million players need roughly 0.5 to 1 GB of RAM, which is one node.
Why not SQL: "top 100" is easy with an index on score, but "my rank" is count(*) WHERE score > mine, which counts up to 5 million index entries for a low-ranked player, 20,000 times a second.
Reads: the top 100 is the same for everyone, so it is read-hot. Cache the rendered list in each app server for one second, and the sorted set sees a few requests a second for it instead of 30,000. Serve rank lookups from the primary or from replicas.
Durability: the match results are the source of truth. If the set is lost, rebuild it from them; turn on AOF so that a restart usually does not need a rebuild. Say the cost: during a rebuild, ranks are missing or partial.
Cheat sheet: pressure β move β cost
| Pressure in the requirements | Move | What you give up |
|---|---|---|
| Datasets with different shapes, rates, lifetimes | Place each by access pattern; relational by default | More systems; copies to keep in sync |
| Frequent query scans | Composite index: equality columns, then range or sort | Slower writes, more storage |
| Race on one row (booking, balance, stock) | Conditional UPDATE, check rows affected |
Contention on that row |
| Rule across rows (one active trip) | Constraint or unique index, or Serializable with retry | Retries; rules stay inside one database |
| Several writes must commit together | One transaction in one database | Cross-service atomicity needs a saga or an outbox |
| Machine may die; money data | Synchronous or quorum standby | Commit latency |
| Reads outgrow the primary | Async replicas plus read-your-writes routing | Stale reads, lag to monitor |
| Same data read far more than written | Cache-aside, delete on write, TTL | Bounded staleness; cold starts |
| Hot key expiring | Single-flight, serve stale, early refresh | Brief staleness |
| Ephemeral, huge write rate, rebuildable | In-memory store, no persistence, keys split so no node is hot | Loss on restart (rebuilt by the next writes) |
| Data or writes past one primary | Shard by the key the queries use, with many logical partitions | Scatter-gather for other queries, resharding work |
| Read-hot key | Cache or replicate it | Staleness |
| Write-hot key | Split into sub-keys, aggregate, or queue | Reads must combine; small loss window if aggregating |
| Human error, corruption | Backups with point-in-time restore, tested | Restore time (RTO) |
Before moving on, explain one design aloud as you would to an interviewer: the ride-hailing app, or a system you know. For each dataset, answer the five questions. Name the store, the index for the busiest query, each invariant and the statement or constraint that enforces it, the replication mode with its RPO, the cache pattern with its staleness bound, and the shard key with its hot-key plan. Then answer the three follow-ups you will get: "What if the cache disappears?", "What if the primary dies mid-commit?" and "What if one city is ten times bigger than the rest?" If you can do that without notes, you are ready for any storage deep dive.
Next: Messaging & Queues for the logs, change-data-capture pipes and write buffers used here; Advanced Topics & Final Prep for consistent hashing, leader election, sagas and the outbox pattern; and URL Shortener (e.g. TinyURL) to apply these moves in one complete design.