Technologies that implement this pattern: PostgreSQL · Cassandra · DynamoDB · Apache Kafka · Time Series DBs · Redis
Why This Matters#
Reads can be copied. Writes cannot. Every read-scaling technique works by serving a stale copy from somewhere cheaper; a write has to land in the authoritative place, become durable, and — usually — be applied to every index, replica and derived view that depends on it. That is why write scaling is where systems hit walls that money alone does not fix: a single primary has one WAL, one fsync path and one lock manager, and adding read replicas makes the write problem slightly worse, not better.
Most candidates treat write scaling as a sharding question: "shard by user ID." Staff engineers treat it as a write-amplification and durability-point question. One logical write from a user typically becomes 5–30 physical writes: the WAL, the heap, three secondary indexes, two replicas, a CDC event, a cache invalidation, a search index update. Before choosing a partition key, a Staff engineer asks: when is this write acknowledged, what does it cost per logical write, and which of those physical writes have to happen before the user gets a 200? Often the cheapest scale-up is deleting an index or moving four of those writes off the synchronous path.
The second reframe: write scaling is a skew problem. Aggregate throughput is easy to buy — add shards. What breaks is the one partition that gets 40% of the traffic: the celebrity's inbox, today's date in a time-keyed table, the largest tenant, the counter everyone increments. Sharding spreads keys, not load on a key.
If you can walk an interviewer from "what's the durability point and amplification factor" to "how the write path is partitioned and what happens to the hottest partition" to "how we reshard without downtime and who owns the migration", you are answering at Staff level.
The 60-Second Version#
- Count physical writes per logical write first. WAL + heap + each index + each sync replica + each derived store. Cutting amplification from 12× to 6× doubles capacity without new hardware. Every secondary index on a hot table costs ~10–30% of insert throughput.
- A single relational primary goes further than people think. A well-tuned Postgres/MySQL primary on NVMe sustains roughly 10–50K small transactional writes/sec with group commit. Below that, batching, fewer indexes and connection pooling beat sharding.
- Batching is the cheapest multiplier. One fsync per 100 rows instead of per row is a 10–50× throughput gain. Kafka producers, multi-row
INSERT, andCOPYall exploit it. The price is latency (5–50ms linger) and a bigger blast radius per failed batch. - Partition on the key that spreads writes and serves the hot read.
user_idusually does both;created_atdoes neither (all writes hit the newest partition). Hash for spread, range for scans; compound keys when you need both. - Budget for the hottest partition, not the average. DynamoDB caps a partition at ~1,000 WCU/s; a Kafka partition handles ~10–50 MB/s; a Cassandra partition beyond ~100 MB degrades. Plan the top 0.1% of keys explicitly — salt, split, or aggregate them.
- Move work off the acknowledged path. Acknowledge after the write hits a durable log (Kafka with
acks=all, a WAL), then materialize indexes, counters and views asynchronously. That turns a 40ms synchronous fan-out into a 5ms append — and makes "how stale is the view?" a metric someone must own.
The Problem#
A metrics product ingests 2M data points per second from 500K agents. A social app records 300K likes per second during a live event, 60% of them on 20 posts. A ride-hailing platform writes a GPS ping every 4 seconds from 1.5M active drivers — ~375K writes/sec of tiny rows that are worthless after ten minutes. A payments ledger must durably record 15K transactions/sec with zero loss and strict per-account ordering. None of these fit on one primary, and they fail in different ways: the metrics pipeline drowns in index maintenance, the likes counter melts one hot row, the GPS stream wastes durable storage on data nobody reads, and the ledger can't afford to acknowledge a write that isn't on two disks. "Shard it" is the start of the answer for all four and the whole answer for none.
Case Studies That Use This Pattern#
- Database Sharding — The flagship: choosing keys, routing, resharding live, and cross-shard queries
- Ad Click Aggregator — Millions of events/sec; pre-aggregate before the store, exactly-once counts for billing
- Metrics & Monitoring — Write-dominated time series; compression, downsampling, cardinality as the real write cost
- Chat Messaging — Per-conversation partitioning with ordering; hot group chats as the skew problem
- Leaderboard — Hot counters; sharded increments and periodic merge
- Ride Hailing & Delivery — High-frequency location writes that belong in memory, not on disk
- Payment Processing — Durable, ordered, append-only ledger writes where losing one is unacceptable
- Web Crawler — Bulk writes of fetched pages and frontier state; batching and backpressure
For the mechanics of partitioning schemes themselves (hash vs range, consistent hashing, rebalancing), see Sharding & Partitioning and Consistent Hashing.
The Four Intents#
"Scale the writes" means at least four different things. They lead to different stores, different acknowledgment points and different failure tolerances.
| Intent | Constraint | Strategy | Failure Mode | Correctness Bar |
|---|---|---|---|---|
| Durable business records (orders, ledger entries, messages) | Acked write must survive a node loss; often per-entity ordering | Partitioned OLTP (sharded Postgres/Vitess, DynamoDB, Spanner) with sync replication | Hot partition; resharding bugs; cross-shard transactions | Zero loss after ack, per-key order |
| High-volume telemetry (metrics, logs, clicks, GPS) | 100K–10M events/sec; individual loss tolerable | Durable log (Kafka) + batched, compressed writes into LSM/TSDB/columnar stores | Consumer lag; compaction debt; cardinality explosion | At-least-once; ≤ 0.01–0.1% loss acceptable if signed off |
| Hot aggregates (counters, likes, quotas, leaderboards) | Huge write rate on few logical keys | Pre-aggregate in memory or in the stream; sharded counters; periodic flush | Lost increments on crash; skew on one key | Bounded error or exact-after-reconcile |
| Burst absorption (flash events, backfills, batch imports) | Peak 10–50× average for minutes | Queue/log as shock absorber; rate-limited consumers; backpressure | Queue growth beyond retention; downstream overload | Same as underlying intent, delayed |
🎯 Staff Move: "Most of this volume is telemetry — location pings at 375K/s that are worthless in ten minutes. I'll treat that as intent two: in-memory latest-position store plus a Kafka log sampled into cold storage. The trip and payment records are intent one and get a small, properly replicated OLTP path. Mixing those two in one database is how we'd end up paying ledger prices for GPS noise."
The Core Tradeoff#
| Strategy | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Vertical scale + tuning (bigger box, group commit, fewer indexes, pooler) | Zero architecture change; often 2–5× headroom | Hard ceiling; one failure domain; replica lag grows with write rate | Future team when the ceiling arrives during a peak |
| Batching / group commit | 10–50× fewer fsyncs and round trips | Adds 5–50ms latency; partial-batch failure semantics | Latency-sensitive callers; error-handling code |
| Hash sharding (by tenant, user, entity) | Near-linear write scale; isolated failure per shard | Cross-shard queries and transactions; resharding is a project; hot keys still hot | App teams (routing, no joins); platform (resharding) |
| Range / time partitioning | Efficient scans, cheap retention (drop old partitions) | All new writes hit the newest range — a moving hot spot | The latest partition's node; on-call during peaks |
| Log-first, materialize later (Kafka → stores) | Fast durable ack; replay; many consumers | Read-your-writes breaks; views lag; two systems to run | Users seeing stale views; team owning consumer lag |
| LSM / write-optimized stores (Cassandra, RocksDB, TSDBs) | Sequential writes, 50–100K+ writes/sec/node | Read amplification; compaction debt stalls writes; tombstones | Read latency; on-call during compaction storms |
| Pre-aggregation (in-memory, stream windows) | Collapses N writes into 1 | Crash loses unflushed deltas; exactness requires checkpointing | Accuracy (bounded error) or complexity (exactly-once) |
Staff Default Position#
Shrink the write before you split it: cut amplification, batch, and move derived writes off the ack path — then partition on a key that spreads load and serves the dominant read, with an explicit plan for the hottest keys.
Start on one well-tuned primary until measured write load passes ~50–60% of its sustained capacity. Remove or defer indexes that don't serve a hot read; batch; acknowledge at the durable log and materialize views asynchronously where product accepts the lag. When you shard, shard by the entity that owns the invariant (tenant, user, account) using hash partitioning with many more logical partitions than physical nodes (e.g., 4,096 logical → 16 physical), so resharding is moving partitions, not rehashing rows. Name the hot-key strategy before launch: salting for append-only keys, pre-aggregation for counters, dedicated shards for whale tenants. Separate telemetry from records — they have different durability bars and should not share a failure domain.
When to Deviate#
- Strict cross-entity transactions at scale — If the business genuinely needs ACID across arbitrary rows (global inventory transfers, cross-account ledger moves), a distributed SQL store (Spanner, CockroachDB, YugabyteDB) buys sharding without giving up transactions, at ~5–20ms extra commit latency in-region and much more cross-region. Cheaper than building 2PC yourself.
- Writes that don't need durability — Live location, presence, typing indicators. Keep them in memory (Redis with no AOF fsync, or an in-process grid) and accept loss on failover. Writing these to a durable store is pure cost.
- Low scale with headroom — Under ~5K writes/sec on a modern primary at < 40% utilization, sharding buys complexity and nothing else. Partition tables for retention, not throughput.
- Bulk/backfill loads — Use the store's bulk path (
COPY, SSTable bulk load, S3 → DynamoDB import) with indexes disabled or built after, not millions of individualINSERTs through the OLTP path.
The L5 → L6 → L7 Contrast#
| Behavior | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| First move | "Shard the database by user ID" | "What's the ack point, how many physical writes per logical write, and what's the hottest key's rate?" | "Which of our write paths are records vs telemetry, and should they share infrastructure, budgets and durability standards at all?" |
| Partitioning | Hash by ID, N = number of nodes | Many logical partitions over few nodes; key chosen for spread and the dominant read; hot-key plan | Sets the org's partitioning standard and owns the resharding platform so teams don't hand-roll migrations |
| Durability | "Postgres is durable" | Explicit ack point (sync replica, acks=all), loss budget for telemetry signed off by product | Defines durability tiers (D0 ledger, D1 records, D2 telemetry) with cost per tier and maps every write path to one |
| Failure | Add retries | Backpressure, bounded queues, compaction and lag alerts, hot-shard detection | Cell architecture so one shard's failure is a known fraction of users; resharding rehearsed as a routine operation |
| Cost | "More shards" | "Dropping two indexes buys 40% headroom — six months before we need to shard" | Prices $/million writes per tier; moves telemetry off premium storage; retires write paths nobody reads |
| Ownership | Service team | Service owns key choice; DB platform owns shard topology | Platform owns routing, resharding and the storage tiers; product teams own keys, loss budgets and retention |
Why "First move" separates levels
Sharding by user ID is often the right destination. Opening with it skips the cheaper moves (fewer indexes, batching, async derivation) that buy 2–5× and defer a multi-quarter migration. It also skips the question that decides whether sharding will even work: is the load skewed? The Staff candidate measures amplification and skew first. The Principal candidate notices that half the org's write volume is telemetry stored at ledger prices and fixes the classification before anyone shards anything.
Why "Durability" separates levels
"The database is durable" hides the actual decision: when is a write acknowledged? After the local WAL fsync? After one synchronous replica? After a Kafka quorum? Each answer has a latency cost (0.1ms, 1–2ms, 5–10ms in-region) and a loss window. Staff names the ack point per write path. Principal turns it into tiers with prices so teams stop defaulting to the most expensive one.
Why "Ownership" separates levels
Resharding incidents come from teams doing a once-in-three-years operation for the first time. Staff makes sure the owning team has a runbook. Principal makes resharding a platform capability exercised monthly, so it is boring by the time anyone needs it.
The Five Fault Lines#
| # | Fault Line | The Tension |
|---|---|---|
| 1 | Synchronous vs Deferred Derivation | Fresh indexes and views vs a fast, small acknowledged write |
| 2 | Hash vs Range Partitioning | Even spread vs efficient scans and retention |
| 3 | Exact vs Aggregated Writes | Every event stored vs N events collapsed into one |
| 4 | Durability Point vs Latency | Ack after more copies vs ack faster with a loss window |
| 5 | Shard Now vs Shard Later | Up-front complexity vs a live migration under pressure |
Fault Line 1: Synchronous vs Deferred Derivation#
A product write that synchronously updates the row, four indexes, a search document, a cache and a feed costs ~30–60ms and fails if any of them fails. Writing the row plus an outbox record in one transaction (~3–5ms) and deriving the rest from CDC makes the acknowledged path small and robust. Who pays: deferred derivation is paid by readers who see a view that is 100ms–5s behind, and by the team that owns consumer lag; synchronous derivation is paid by write latency, availability (the product of every dependency's uptime), and capacity. Staff default: keep in the transaction only what enforces an invariant or serves read-your-writes for the writer; derive everything else from the log. Deviate when: the derived structure is the invariant (a unique index on email) — that stays synchronous.
Fault Line 2: Hash vs Range Partitioning#
Hash partitioning on user_id spreads writes evenly but turns "all orders from yesterday" into a scatter across every shard. Range partitioning on (tenant, created_at) makes time scans and retention cheap (drop a partition, no delete storm) but sends all of today's writes to the newest range. Who pays: hash — analytical and cross-entity queries (scatter-gather, or a separate OLAP copy); range — the node owning the head of the range, which runs hot while others idle. Staff default: hash on the owning entity for OLTP; time-partition within each shard for retention; analytics from a CDC-fed warehouse. Deviate when: the workload is a time series with append-only writes and time-bounded reads — use compound keys like (metric_id, hour_bucket) so writes spread across series while scans stay local.
The picture is the argument: a pure time key concentrates every write on one node, hash spreads them, and a compound key gets spread for writes and locality for time-bounded scans.
Fault Line 3: Exact vs Aggregated Writes#
A like button at 300K/s could write 300K rows per second, or a stream processor could collapse them into per-post deltas every second and write ~20K counter updates. Telemetry can be downsampled at the edge (agents send 10s summaries instead of 1s samples). Who pays: aggregation loses detail (can't answer "who liked it at 12:03:07?" from the counter) and risks losing unflushed deltas on crash; exact writes pay storage and a write path 10–100× larger. Staff default: store events raw in a cheap log tier (Kafka → object storage) for audit and recompute; serve counters from pre-aggregated values; make the aggregation exactly-once only when money depends on it. Deviate when: individual events are the product (messages, transactions) — never aggregate those.
Fault Line 4: Durability Point vs Latency#
| Ack point | Added latency (in-region) | Loss window on failure |
|---|---|---|
| Memory only | ~0 | Everything since last flush |
| Local WAL fsync | 0.1–1ms (NVMe) | Node disk loss |
| Sync replica (1 of N) | +1–2ms | Simultaneous loss of two nodes |
Quorum log (Kafka acks=all, min.insync.replicas=2) | +2–10ms | Correlated loss of ISR |
| Cross-region sync | +30–150ms | Region loss covered |
Who pays: a later ack point is paid by every writer's latency; an earlier one is paid by the users whose data vanishes in the rare failure — and by whoever explains it. Staff default: quorum or sync-replica ack for records; local or batched ack for telemetry; cross-region sync only for the small set of writes where region loss would be a regulatory event. Deviate when: product explicitly accepts RPO > 0 for a record type (e.g., draft autosaves) — then async is fine and should be documented.
Fault Line 5: Shard Now vs Shard Later#
Sharding on day one costs routing code, no cross-entity joins, more operational surface — for load that may never come. Sharding later means a live migration of a hot database under growth, typically 2–4 quarters with dual writes, backfills, verification and cutover. Who pays: early — every feature team, in slower development; late — the platform and on-call, in a high-risk migration during the period of fastest growth. Staff default: don't shard early, but choose the shard key early: make every table carry the future shard key (tenant_id/user_id), avoid cross-tenant foreign keys and sequences, and use IDs that don't depend on a single sequence. That makes a later split a mechanical operation. Deviate when: the growth curve is known (B2B onboarding a 10× customer next quarter) — shard ahead of the step.
Common Interview Mistakes#
| What Candidates Say | What Interviewers Hear | What Staff Engineers Say |
|---|---|---|
| "Shard by user ID" (first sentence) | "Jumped to the expensive answer" | "First: 9 indexes on this table, 4 unused. Dropping them and batching buys ~2× — then shard by user_id with 4,096 logical partitions." |
| "Add read replicas to handle the write load" | "Doesn't understand replication" | "Replicas replay every write; they add write work, not capacity. Writes scale by partitioning or by doing fewer of them." |
| "Partition by timestamp for easy cleanup" | "Hasn't seen the hot-head problem" | "Time partitions inside each hash shard: retention is a partition drop, writes stay spread." |
| "Kafka makes it scalable" | "Moved the problem, didn't solve it" | "Kafka makes the ack cheap and absorbs bursts. The consumer still has to write at the sustained rate, and lag is the new staleness metric." |
| "Use consistent hashing" (as the whole answer) | "Knows the term, not the operations" | "Fixed logical partitions mapped to nodes; moving a partition is a copy plus a routing flip, and we rehearse it monthly." |
| "Counters just increment in the DB" | "Hot row at 300K/s" | "Aggregate per post per second in the stream, write deltas; raw events go to the log for audit." |
Quick Reference#
Staff Sentence Templates#
"Each logical write here becomes about [N] physical writes: [WAL, heap, K indexes, R replicas, CDC, cache]. Before sharding, I'd cut that to [M] by [dropping / deferring] [which ones] — that alone is [N/M]× headroom."
"I'll acknowledge at [sync replica / Kafka acks=all / local WAL], which costs about [X ms] and loses at most [window] on [failure]. For [telemetry path], product signed off on losing up to [0.1%], so it acks [earlier]."
"I'll shard by [entity] because it owns the invariant and serves the dominant read [query]. The hottest [entity] does [R] writes/sec against a partition ceiling of [C], so [it fits / I'll salt it into K sub-keys]."
"Resharding is [moving logical partitions / a dual-write migration]. I'd make it routine — [platform team] runs it monthly — so it's not a first-time operation during our biggest growth quarter."
Implementation Deep Dive#
1. Shrink the Write — Outbox + CDC Instead of Synchronous Fan-out (PostgreSQL)#
The cheapest write scaling is doing less inside the transaction. Keep the invariant-bearing write; emit an event in the same transaction; derive everything else.
BEGIN;
INSERT INTO orders (order_id, user_id, total, status, created_at)
VALUES ($1, $2, $3, 'placed', now()); -- PK + 1 index (user_id, created_at)
INSERT INTO outbox (id, aggregate_id, type, payload)
VALUES (gen_random_uuid(), $1, 'OrderPlaced', $json); -- tiny, append-only
COMMIT; -- ~2-4 ms with sync replica
-- Debezium tails the WAL -> Kafka topic orders.events (partitioned by user_id)
-- Consumers, each owning its lag:
-- search-indexer -> Elasticsearch (lag SLO 5s)
-- analytics-sink -> warehouse, hourly files (lag SLO 1h)
-- user-stats -> per-user aggregates (lag SLO 30s)
-- cache-invalidator -> Redis DEL (lag SLO 1s)
Numbers: moving four synchronous side effects (~8–15ms each, each a failure point) into async consumers typically takes p99 write latency from ~60ms to ~5ms and raises the primary's sustainable write rate 2–3× because indexes for search/analytics no longer live on it. The cost: the order list in search may be up to 5s stale, so the "my orders" page for the writer reads from the primary for 10s after a write.
🎯 Staff Insight: The outbox is not just a reliability pattern — it's a write-scaling pattern. Every index you delete from the OLTP primary because a downstream system serves that query is capacity you didn't have to shard for.
2. Logical Partitions and Routing — Sharded Postgres / Vitess-style#
LOGICAL_PARTITIONS = 4096 # fixed forever; nodes change, this doesn't
function partition_of(user_id):
return murmur3(user_id) % LOGICAL_PARTITIONS
# Routing table, stored in etcd, cached in every app/proxy with a watch
# partitions 0-255 -> shard-01 (primary pg-01a, replica pg-01b)
# partitions 256-511 -> shard-02
# ...
function route(user_id):
p = partition_of(user_id)
return routing_table.lookup(p) # version-stamped
# Every row carries user_id; every table is co-partitioned by it.
# No cross-partition foreign keys. IDs embed the partition for routing without lookup:
# order_id = (timestamp_ms << 22) | (partition << 10) | seq (Snowflake-style)
Why 4,096: with 16 physical shards each holds 256 logical partitions; growing to 32 shards means moving 128 partitions per old shard — a copy and a routing-table flip each, not a rehash of every row. Instagram's published scheme (thousands of logical shards mapped onto far fewer Postgres servers, with the shard ID embedded in every ID) is the canonical public example.
3. Hot-Key Write Salting — DynamoDB / Cassandra#
For append-heavy keys that exceed one partition's ceiling (a viral post's comments, a global event stream, today's bucket):
SALT_BUCKETS = {"default": 1, "hot": 16} # promoted by a hot-key detector
function write_comment(post_id, comment):
n = salt_buckets_for(post_id) # 1 normally, 16 when hot
pk = post_id + "#" + (hash(comment.id) % n)
ddb.PutItem(pk=pk, sk=comment.ts + "#" + comment.id, ...)
function read_recent(post_id, limit):
n = salt_buckets_for(post_id)
pages = parallel([ddb.Query(pk=post_id + "#" + i, desc, limit) for i in 0..n-1])
return merge_by_sk(pages).take(limit) # n-way merge
Numbers: DynamoDB allows ~1,000 WCU/s and ~3,000 RCU/s per partition; 16 salt buckets lift a hot key to ~16K writes/s. The read side pays a 16-way fan-out — acceptable because hot posts are also the most cached. The salt count must be discoverable by readers (a small config item, cached), and lowering it requires a migration — it's a one-way-ish door per key.
4. Pre-Aggregation with Exactly-Once Flush — Kafka + Stream Processor#
# Likes arrive at 300K/s on topic likes (key = post_id)
stream = kafka.source("likes", group="like-counter")
counts = stream
.key_by(post_id)
.window(tumbling(1s))
.aggregate(count) # ~20K distinct posts/sec
counts.sink(postgres, mode="idempotent upsert"):
INSERT INTO post_likes (post_id, window_start, delta)
VALUES ($1, $2, $3)
ON CONFLICT (post_id, window_start) DO NOTHING; -- replay-safe
# Totals: periodic roll-up or materialized view
# likes_total(post_id) = sum(delta) — served from cache, refreshed every 1-2s
Why it works: the database sees ~20K small upserts per second instead of 300K row inserts or 300K updates to 20 hot rows. Checkpointing the stream offsets together with the windowed state means a crash replays from the last checkpoint and the (post_id, window_start) key makes the replayed writes no-ops. See Flink & Stream Processing for the checkpoint mechanics.
Technique Comparison
| Technique | Throughput Gain | Latency Cost | Consistency Cost | Operational Burden |
|---|---|---|---|---|
| Drop/defer indexes | 1.3–3× | None (faster) | Derived views lag | Low |
| Batching / group commit | 10–50× fsync reduction | +5–50ms | Batch-level failure | Low |
| Outbox + CDC | 2–3× on primary | Lower ack latency | View staleness | Medium (consumer lag) |
| Hash sharding | ~N× (N shards) | +routing hop (~0.1ms) | No cross-shard ACID | High (resharding) |
| Hot-key salting | ×salt count per key | Fan-out on read | Merge ordering | Medium |
| Pre-aggregation | 10–100× fewer writes | Window delay (1–60s) | Detail lost; exactness needs checkpoints | Medium–High |
| LSM / TSDB store | 3–10× vs B-tree for inserts | Read amplification | Compaction stalls | Medium |
Architecture Diagram#
How to narrate it: two write tiers that never share a failure domain. Records go through a router to partitioned primaries with synchronous replication and nothing else in the transaction; everything derived comes off the WAL. Telemetry never touches an OLTP database: it is batched into a log with a cheaper ack, aggregated, and stored in engines built for appends. The arrows into Kafka are where "write latency" stops and "view freshness" begins — and each consumer has its own lag SLO and owner.
Failure Scenarios#
1. The Moving Hot Head — Time-Keyed Table Saturates One Node#
An events table in Cassandra is keyed by (event_date) with clustering on timestamp. It worked in testing at 2K/s. At launch, writes reach 60K/s.
t=0 Midnight UTC: all writes move to partition '2026-03-14' on 3 replicas.
t=+2h Partition grows to 4 GB; node CPU 95% on 3 of 24 nodes; others at 10%.
t=+3h Write timeouts 8%; coordinators retry; hinted handoffs pile up.
t=+5h Compaction on the hot replicas falls behind; read latency for today 3s.
t=+6h On-call adds nodes. No effect: one partition lives on 3 replicas regardless.
t=+9h Emergency deploy: key changed to (event_date, bucket 0-63). Backfill read path.
Detection: node.write_latency skew > 5× median; table.max_partition_bytes > 100 MB; coordinator.write_timeouts rising while cluster-average CPU is low.
Blast radius: all writes to the table, plus reads for "today" — the most-read data.
Mitigation: add a bucket component to the key; raise client timeouts temporarily; throttle non-critical producers.
Prevention: design review rule — no partition key may be a pure time value; capacity test at 10× expected with production-like key distribution; alert on max partition size.
Owner: owning service team for the data model; data platform for the review rule and partition-size alerting.
🎯 Staff Insight: "Add nodes" did nothing because the bottleneck was a single key, not the cluster. If adding capacity doesn't move the metric, the problem is skew — look at the per-key distribution before buying hardware.
2. Resharding Dual-Write Gap — 0.3% of Orders Missing on the New Shards#
A team splits 8 shards into 16. The plan: dual-write to old and new, backfill historical rows, verify, cut reads over, stop writing old. The dual-write was implemented in application code, not from the WAL.
t=0 Dual-write enabled in app. Backfill job starts copying historical rows.
t=+2d Backfill complete. Spot-check verification passes (1,000 random rows).
t=+3d Reads cut over to new shards for 10 percent of users.
t=+3d+4h Support tickets: "my order disappeared". 0.3 percent of recent orders absent.
t=+3d+5h Root cause: an admin bulk-edit tool and a refund batch job write directly
to the old shards. Neither went through the dual-write code path.
t=+3d+6h Reads rolled back to old shards. Diff job finds 41K divergent rows.
Detection: migration.row_diff_count from a full keyed checksum (not sampling); orders.not_found_after_create rate.
Blast radius: 10% cohort for 4 hours; reconciliation work for every divergent row.
Mitigation: roll back reads (old shards were still authoritative — the reason you keep them writable until verified).
Prevention: copy from the WAL/CDC stream, not from application dual-writes, so every write path is captured; verify with full per-partition checksums; cut over per logical partition with automatic rollback on diff.
Owner: DB platform owns the migration tooling; the product team owns the inventory of every write path to its tables.
3. Compaction Debt — LSM Store Stalls Writes During a Backfill#
A metrics team backfills 30 days of data into a RocksDB-based TSDB at 5× normal ingest while live writes continue.
t=0 Backfill starts at 1M points/s on top of 2M live.
t=+40min L0 file count exceeds slowdown threshold; writes throttled 50 percent.
t=+55min L0 hits stop threshold. Writes stall entirely on 6 of 20 nodes.
t=+56min Agents buffer locally (5 min capacity), then drop. Dashboards show gaps.
t=+70min Backfill paused. Compaction catches up in 25 minutes.
t=+95min Live ingest normal. 14 minutes of metrics lost for 30 percent of hosts.
Detection: lsm.l0_files vs slowdown/stop thresholds; lsm.pending_compaction_bytes; ingest.write_stall_seconds; agent.buffer_drop_count.
Blast radius: live observability for a third of the fleet — during which nobody could see other problems.
Mitigation: pause or rate-limit the backfill; shed lowest-priority series.
Prevention: backfills go through a separate, rate-limited path (or bulk SST ingestion that bypasses L0); backfill rate is a fraction of measured compaction headroom; live traffic has priority via separate queues.
Owner: metrics platform team; backfill requests go through them.
Operational Reality Matrix#
| Failure | Detection Signal | Blast Radius | Mitigation | Owner |
|---|---|---|---|---|
| Hot partition / key | partition.write_rate top-K vs ceiling; node skew > 3× | All writes on that partition | Salt, split, dedicated shard | Service (key design) + platform (detection) |
| Primary at write capacity | db.wal_bytes_per_sec, db.commit_latency p99 > 10ms | Whole shard | Batch, shed async writers, defer indexes | Service + DB platform |
| Replica lag from write burst | db.replica.lag_seconds > 5s | Read-after-write, failover RPO | Throttle bulk writers, separate bulk path | DB platform |
| Consumer lag on derived stores | kafka.consumer_lag_seconds > SLO | Stale search, analytics, caches | Scale consumers, prioritize topics | Each consumer's owning team |
| Compaction debt | lsm.pending_compaction_bytes, write stalls | Ingest for affected nodes | Rate-limit backfills, add IO | Storage/metrics platform |
| Resharding divergence | migration.row_diff_count > 0 | Users in migrated cohort | Roll back per partition | DB platform |
| Unbounded cardinality | tsdb.active_series growth > 20%/day | Whole TSDB cluster | Drop/aggregate offending label | Metrics platform + emitting team |
The Principal Lens#
Why L7 Sees This Problem Differently#
A Staff engineer scales the write path of one system. A Principal engineer looks at the org's total write volume and finds that 70% of it is telemetry and derived data stored on infrastructure priced for business records, that three teams are independently planning resharding projects on homegrown routing layers, and that nobody can say which writes the company would actually miss after a region failure. At org scale, write scaling becomes a classification and tiering problem: which writes are records, which are telemetry, which are derived and can be rebuilt — and what durability, retention and cost each tier gets. The highest-leverage move is often not making writes faster, but deciding which ones don't need to exist, don't need to be durable, or don't need to live more than seven days.
The Org-Level Fault Line#
One storage platform with durability tiers vs each team choosing and operating its own write path.
| Option | What Works | What Breaks | Who Pays |
|---|---|---|---|
| Every team picks its store | Fit-for-purpose, fast starts | 6 database engines, 3 sharding schemes, nobody expert in any; resharding learned per team | On-call fatigue; security/compliance reviews ×6 |
| One mandated database for all writes | Deep expertise, single toolchain | Telemetry crushes the OLTP fleet; records priced like telemetry or vice versa | Everyone during the noisy-neighbor incident |
| Tiered platform: records / events / telemetry, each with a paved store, routing and resharding | Teams declare a tier; platform owns operations; costs visible per tier | Platform must support 3–4 engines well; exceptions needed for genuinely novel workloads | Platform headcount (6–10 across tiers) |
The Principal default is tiering: product teams declare durability (D0–D2), retention and access pattern; the platform maps that to a paved store and owns routing, resharding and backups.
Cost Model#
Assumptions: OLTP storage on replicated NVMe ~$0.30–0.50/GB-month all-in (3 copies), log storage ~$0.10/GB-month, object storage ~$0.02/GB-month, engineer ~$25K/month. "Writes" means logical writes at peak.
| Scale | Write Rate (records / telemetry) | Infra | People / On-call | Rough Monthly Total |
|---|---|---|---|---|
| Startup | 1K / 10K per sec | 1 Postgres primary + replica ( | 0.3 FTE; team on-call | ~$4K + ~$8K people |
| Growth | 15K / 500K per sec | 16-shard Postgres ( | 2 FTE DB platform; shared rotation | ~$85K + ~$50K people |
| Large | 200K / 10M per sec | Distributed SQL or 128 shards ( | 8–10 engineers across storage platforms; dedicated rotations | ~$650K + ~$250K people |
The Principal observation: at growth and large scale, the biggest lever is tier placement, not engine choice. Moving telemetry that was landing in the OLTP fleet (with 3× replication on NVMe and 1-year retention) to a log + object-store tier with 7-day hot retention commonly cuts that data's storage cost by 10–20×. The second lever is retention: half of most companies' write volume is never read after 30 days.
The 3-Year Evolution Path#
One-Way Doors vs Two-Way Doors#
| Decision | Door Type | Reversibility Cost |
|---|---|---|
| Batch size, linger, index set | Two-way | Config or migration in hours |
| Number of physical shards (with logical partitions) | Two-way | Move partitions; routine |
| Shard key choice | One-way | Re-keying every table and every query; 2–4 quarters |
| Number of logical partitions | One-way-ish | Changing it is a full rehash; pick generously (1K–16K) |
| ID format without embedded routing info | One-way-ish | Every lookup needs a directory; hard to retrofit |
| Exposing "write acknowledged = visible everywhere" in a public API | One-way | Clients depend on synchronous derivation forever |
| Telemetry vs record tier for a data type | Two-way (if events are logged raw) | Re-derive from the log |
The Standard I'd Write#
RFC: Write Path and Durability Tiering Standard (v1)
Scope: Every service that persists data, and every new table, topic or bucket.
MUST:
- Each write path declares a durability tier: D0 (cross-region sync; ledgers, regulatory), D1 (in-region quorum; business records), D2 (best-effort, batched; telemetry), plus retention and an owner.
- Every OLTP table carries the owning entity's partition key; no cross-partition foreign keys or global sequences.
- Derived data (search, analytics, caches, counters) MUST be produced from CDC/outbox, not from synchronous dual writes in request handlers.
- Partition keys MUST NOT be a bare timestamp or low-cardinality value; designs include the expected top-10 key rates and the hot-key plan.
- Backfills and bulk loads use the platform bulk path with rate limits set by the platform.
SHOULD: Batch writes ≥ 10 rows where latency allows; review indexes quarterly against query logs; retain raw telemetry hot ≤ 7 days by default.
Exceptions: Storage platform review; time-boxed to two quarters.
Success metrics: storage cost per million writes per tier down 25% YoY; zero Sev-1s from resharding; no OLTP fleet hosting D2 data; median resharding operation < 1 week.
What I'd Tell the VP#
Our databases are approaching the point where we either split them or slow down product work, and splitting is a year-long project if we start it in a panic. Most of what we write today is device and activity data that we keep on our most expensive storage and rarely read. I'm proposing two things: move that data to cheaper storage built for it, which buys us roughly a year of headroom and saves around a third of our database spend; and build the splitting capability once, as a platform, so that when the core databases do need to grow, it's a routine operation instead of a risky migration. It needs about four engineers for three quarters.
Principal Interview Signals#
| Signal | What It Sounds Like |
|---|---|
| Classification before scaling | "Seventy percent of these writes are telemetry. They shouldn't be on the same tier — or the same budget — as orders." |
| Pricing the tier | "Moving location history from replicated NVMe to log plus object storage is about 15× cheaper per GB, and nobody reads it after a week." |
| Resharding as a capability | "I'd rather run resharding monthly on small partitions than once every three years on the whole fleet." |
| One-way door awareness | "The shard key is the decision I can't take back, so it's the one I want product and data to sign off on." |
| Retention as write scaling | "The cheapest write to scale is the one we delete after seven days." |
Staff answers that L7 interviewers find insufficient:
- "I'd shard this by user_id with 4,096 logical partitions." — Right for the service; doesn't ask whether the other teams building routers should share one.
- "We'll move derived writes to CDC." — Correct, but no standard that stops the next team from adding synchronous dual writes.
- "Telemetry goes to Kafka and a TSDB." — Doesn't price it, set retention defaults, or define who owns the tier.
In the Wild#
Instagram: Logical Shards and IDs That Route Themselves#
Instagram's engineering blog described sharding Postgres into thousands of logical shards (implemented as schemas) mapped onto a much smaller number of physical servers, and generating 64-bit IDs that embed a timestamp, the logical shard ID and a per-shard sequence, generated inside Postgres. Moving logical shards between servers rebalances load without rewriting IDs.
Staff insight: This is the "many logical partitions over few nodes" default with a public reference behind it, plus the detail that makes it work: the routing key is inside the ID, so most lookups need no directory.
Discord: From Cassandra to ScyllaDB for Message Writes#
Discord has written publicly about storing billions of messages partitioned by channel and time bucket, the operational pain they hit with Cassandra (hot partitions from very active channels, garbage-collection pauses, compaction work), and their later migration to ScyllaDB fronted by Rust "data services" that coalesce concurrent requests for the same channel.
Staff insight: The partition key (channel_id, bucket) is the compound-key fix for a time-ordered append workload, and the hot-channel problem is the skew lesson: aggregate capacity was never the issue; the busiest partitions were.
Slack: Vitess for Horizontal MySQL#
Slack has described migrating its MySQL fleet to Vitess, a sharding middleware originally built at YouTube, moving from a workspace-based sharding scheme toward more flexible keyspaces and resharding managed by the platform rather than application code.
Staff insight: The story is ownership: resharding and routing moved from per-team application logic into a platform layer. That's the Principal move this pattern points toward — make the hard, rare operation routine and centralized.
Practice Drill#
Prompt: "We're an IoT company. 800K devices send a reading every 5 seconds (160K writes/sec), plus customers create alerts and dashboards (~200 writes/sec). Everything is in one Postgres cluster that's at 85% CPU, and replica lag hits 30 seconds at peak. Fix it."
Staff Answer
The two workloads have nothing in common and shouldn't share a database. Classify: device readings are telemetry (D2 — product signs off on losing ≤ 0.1% during failures, hot retention 14 days, downsampled after); alerts/dashboards are records (D1). Telemetry path: devices → ingestion gateway → Kafka (acks=1, 20ms linger, partitioned by device_id, ~64 partitions at ~3 MB/s each). Consumers write to a TSDB/columnar store in batches of 5–10K points with compression; a stream job computes 1-minute rollups for dashboards so queries never scan raw points; raw data also lands in object storage as hourly Parquet for 1-year retention at ~1/15th the cost. Alerting: evaluated in the stream against rollups, not by polling the database. Records path: stays in Postgres — at 200 writes/sec it needs no sharding; the 85% CPU and 30s lag were caused by 160K/s of telemetry inserts and their indexes. Migration: dual-publish from the gateway to Kafka, backfill the TSDB from Postgres in rate-limited chunks, cut dashboards over per customer, then stop telemetry inserts and drop those tables (partitioned by day, so a cheap drop). Metrics: kafka.consumer_lag_seconds (SLO 30s), tsdb.ingest_rate, tsdb.active_series, pg.cpu, pg.replica_lag_seconds.
Why this is L6:
- Diagnoses the problem as mixed tiers, not insufficient sharding — the record workload is tiny.
- Sets explicit durability and retention per tier with product sign-off.
- Collapses writes (batching, rollups) and moves alerting off the database, with a safe, reversible migration plan.
What L7 adds:
- Prices it: telemetry on replicated NVMe at 14 days vs log + TSDB + object storage — roughly 10× cheaper per GB — and makes that the default for every telemetry producer, not just this one.
- Sets the org standard that no D2 data lands in OLTP clusters, with a review gate for new tables.
- Notices that per-customer data residency is coming and partitions the telemetry topics by region now — a cheap two-way door that avoids a one-way one later.
Staff Interview Application#
How to Introduce This Pattern#
"Before I shard anything I want to know three things: how many physical writes each logical write turns into, when we acknowledge it, and how skewed the keys are. Usually there's 2× sitting in indexes and synchronous side effects we can move off the request path. Then I'll partition on the entity that owns the invariant, and I'll have a plan for the hottest keys before launch."
Lead with amplification and the ack point, then partitioning, then skew, then how resharding works and who owns it.
When NOT to Use This Pattern#
- Read-dominated workloads: If writes are < 5% of load, scale reads — see Scaling Reads. Sharding for a write problem you don't have makes every read harder.
- Moderate volume on a healthy primary: Under ~5K writes/sec at < 40% utilization, tune and wait. Choose the future shard key now; don't implement sharding.
- Contention on a few keys: If total volume is fine but a handful of keys are hot, that's Dealing with Contention — bucketing or single-writer per key — not fleet-wide sharding.
- Ephemeral state: Presence, live location, typing indicators — keep in memory; durable write scaling is the wrong tool.
Follow-Up Questions to Anticipate#
| Interviewer Asks | What They Are Testing | How to Respond |
|---|---|---|
| "Why not just add replicas?" | Replication fundamentals | "Replicas replay every write. They add read capacity and failover, not write capacity." |
| "How do you pick the shard key?" | Key design | "The entity that owns the invariant and appears in the dominant query. Hash for spread; never a bare timestamp." |
| "How do you reshard without downtime?" | Migration maturity | "Logical partitions over physical nodes; copy a partition from snapshot + WAL, verify by checksum, flip routing, keep the source until verified." |
| "What about the celebrity/whale?" | Skew | "Salt the key into N sub-partitions with a fan-out read, or give the whale a dedicated shard. Plan it before launch." |
| "Cross-shard transactions?" | Consistency tradeoffs | "Design so invariants live within one partition. Where they can't, a saga with compensations — or distributed SQL if it's a large share of writes." |
| "How do you count likes at 300K/s?" | Aggregation | "Window per post per second in the stream, idempotent upsert of deltas; raw events to the log for audit." |
Evaluation Rubric#
| Dimension | Senior (L5) | Staff (L6) | Principal (L7) |
|---|---|---|---|
| Framing | "Shard it" | Amplification, ack point, skew before partitioning | Classifies org write volume into tiers with cost and durability contracts |
| Partitioning | Hash by ID | Logical partitions, key chosen for spread + reads, hot-key plan | Platform-owned routing and resharding as a routine operation |
| Durability | Implicit | Ack point per path, loss budget signed off | D0/D1/D2 tiers, mapped and enforced |
| Failure | Retries | Backpressure, lag SLOs, compaction and skew alerts, reversible migrations | Cells, rehearsed resharding, backfill governance |
| Cost | More nodes | Headroom from fewer indexes and batching | $/million writes per tier; retention as a cost lever |
Strong Hire Signals
| Signal | What It Sounds Like |
|---|---|
| Amplification awareness | "Each order is about 12 physical writes; I can make it 5." |
| Ack-point precision | "Ack at the sync replica, 2ms, survives one node loss." |
| Skew-first partitioning | "What's the top key's write rate against the partition ceiling?" |
| Safe migration | "Copy from the WAL, verify by checksum, cut over per partition." |
Lean No-Hire Signals
| Signal | Why It Misses the Bar |
|---|---|
| Replicas as write scaling | Fundamental misunderstanding of replication |
| Timestamp partition key for high-volume appends | Guarantees a hot head |
| No resharding story | Ships a design that can't grow without an outage |
Common False Positives: Explaining LSM-tree compaction in depth ≠ choosing the right tier. Naming consistent hashing ≠ having a resharding plan. "Use Kafka" ≠ knowing what the consumer must sustain.
Capacity Planning Quick Reference#
Sizing the Write Path#
physical_writes = logical_writes × amplification # WAL + heap + indexes + replicas + derived
shards_needed = ceil(peak_logical_writes / (per_shard_sustained × 0.5)) # 50% headroom
logical_partitions = shards_needed × 64 to 256 # room to grow 10x+ by moving, not rehashing
kafka_partitions = ceil(peak_MBps / 10) # ~10 MB/s per partition, conservative
hot_key_salt = ceil(top_key_writes / (partition_ceiling × 0.5))
storage_per_day = writes_per_sec × bytes_per_write × 86400 × replication / compression
Key Numbers Worth Memorizing#
| Number | Context |
|---|---|
| 10–50K/s | Small transactional writes on one tuned Postgres/MySQL primary (NVMe, group commit) |
| 0.1–1 ms | NVMe fsync; the floor for a durable commit |
| +1–2 ms | Adding a synchronous in-region replica |
| ~1,000 WCU/s | DynamoDB per-partition write ceiling |
| 10–50 MB/s | Practical throughput per Kafka partition |
| 50–100K/s | Writes per node on Cassandra/Scylla-class LSM stores |
| < 100 MB | Healthy Cassandra partition size |
| 10–30% | Insert throughput lost per secondary index on a hot table |
| 10–50× | Gain from batching (fewer fsyncs/round trips) |
| 1K–16K | Logical partition count to fix forever |
Common Pitfalls Checklist#
- Physical writes per logical write counted; unused indexes removed
- Ack point and loss budget declared per write path
- Derived data produced from CDC/outbox, not request-time dual writes
- Partition key is not a bare timestamp; top-key rates measured against ceilings
- Logical partitions ≫ physical nodes; routing table versioned
- Resharding copies from the WAL and verifies with full checksums
- Telemetry and records do not share a failure domain
- Backfills use a rate-limited bulk path below compaction headroom