Hiring BarSupport

Scaling Writes — Cross-Cutting Pattern

Pattern37 min read5 diagrams

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, and COPY all 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_id usually does both; created_at does 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.

IntentConstraintStrategyFailure ModeCorrectness Bar
Durable business records (orders, ledger entries, messages)Acked write must survive a node loss; often per-entity orderingPartitioned OLTP (sharded Postgres/Vitess, DynamoDB, Spanner) with sync replicationHot partition; resharding bugs; cross-shard transactionsZero loss after ack, per-key order
High-volume telemetry (metrics, logs, clicks, GPS)100K–10M events/sec; individual loss tolerableDurable log (Kafka) + batched, compressed writes into LSM/TSDB/columnar storesConsumer lag; compaction debt; cardinality explosionAt-least-once; ≤ 0.01–0.1% loss acceptable if signed off
Hot aggregates (counters, likes, quotas, leaderboards)Huge write rate on few logical keysPre-aggregate in memory or in the stream; sharded counters; periodic flushLost increments on crash; skew on one keyBounded error or exact-after-reconcile
Burst absorption (flash events, backfills, batch imports)Peak 10–50× average for minutesQueue/log as shock absorber; rate-limited consumers; backpressureQueue growth beyond retention; downstream overloadSame 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#

StrategyWhat WorksWhat BreaksWho Pays
Vertical scale + tuning (bigger box, group commit, fewer indexes, pooler)Zero architecture change; often 2–5× headroomHard ceiling; one failure domain; replica lag grows with write rateFuture team when the ceiling arrives during a peak
Batching / group commit10–50× fewer fsyncs and round tripsAdds 5–50ms latency; partial-batch failure semanticsLatency-sensitive callers; error-handling code
Hash sharding (by tenant, user, entity)Near-linear write scale; isolated failure per shardCross-shard queries and transactions; resharding is a project; hot keys still hotApp teams (routing, no joins); platform (resharding)
Range / time partitioningEfficient scans, cheap retention (drop old partitions)All new writes hit the newest range — a moving hot spotThe latest partition's node; on-call during peaks
Log-first, materialize later (Kafka → stores)Fast durable ack; replay; many consumersRead-your-writes breaks; views lag; two systems to runUsers seeing stale views; team owning consumer lag
LSM / write-optimized stores (Cassandra, RocksDB, TSDBs)Sequential writes, 50–100K+ writes/sec/nodeRead amplification; compaction debt stalls writes; tombstonesRead latency; on-call during compaction storms
Pre-aggregation (in-memory, stream windows)Collapses N writes into 1Crash loses unflushed deltas; exactness requires checkpointingAccuracy (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 individual INSERTs through the OLTP path.

The L5 → L6 → L7 Contrast#

BehaviorSenior (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?"
PartitioningHash by ID, N = number of nodesMany logical partitions over few nodes; key chosen for spread and the dominant read; hot-key planSets 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 productDefines durability tiers (D0 ledger, D1 records, D2 telemetry) with cost per tier and maps every write path to one
FailureAdd retriesBackpressure, bounded queues, compaction and lag alerts, hot-shard detectionCell 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
OwnershipService teamService owns key choice; DB platform owns shard topologyPlatform 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 LineThe Tension
1Synchronous vs Deferred DerivationFresh indexes and views vs a fast, small acknowledged write
2Hash vs Range PartitioningEven spread vs efficient scans and retention
3Exact vs Aggregated WritesEvery event stored vs N events collapsed into one
4Durability Point vs LatencyAck after more copies vs ack faster with a loss window
5Shard Now vs Shard LaterUp-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.

Diagram: Fault Line 2: Hash vs Range Partitioning

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 pointAdded latency (in-region)Loss window on failure
Memory only~0Everything since last flush
Local WAL fsync0.1–1ms (NVMe)Node disk loss
Sync replica (1 of N)+1–2msSimultaneous loss of two nodes
Quorum log (Kafka acks=all, min.insync.replicas=2)+2–10msCorrelated loss of ISR
Cross-region sync+30–150msRegion 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 SayWhat Interviewers HearWhat 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#

Diagram: 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

TechniqueThroughput GainLatency CostConsistency CostOperational Burden
Drop/defer indexes1.3–3×None (faster)Derived views lagLow
Batching / group commit10–50× fsync reduction+5–50msBatch-level failureLow
Outbox + CDC2–3× on primaryLower ack latencyView stalenessMedium (consumer lag)
Hash sharding~N× (N shards)+routing hop (~0.1ms)No cross-shard ACIDHigh (resharding)
Hot-key salting×salt count per keyFan-out on readMerge orderingMedium
Pre-aggregation10–100× fewer writesWindow delay (1–60s)Detail lost; exactness needs checkpointsMedium–High
LSM / TSDB store3–10× vs B-tree for insertsRead amplificationCompaction stallsMedium

Architecture Diagram#

Diagram: 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.

Diagram: Architecture Diagram

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#

FailureDetection SignalBlast RadiusMitigationOwner
Hot partition / keypartition.write_rate top-K vs ceiling; node skew > 3×All writes on that partitionSalt, split, dedicated shardService (key design) + platform (detection)
Primary at write capacitydb.wal_bytes_per_sec, db.commit_latency p99 > 10msWhole shardBatch, shed async writers, defer indexesService + DB platform
Replica lag from write burstdb.replica.lag_seconds > 5sRead-after-write, failover RPOThrottle bulk writers, separate bulk pathDB platform
Consumer lag on derived storeskafka.consumer_lag_seconds > SLOStale search, analytics, cachesScale consumers, prioritize topicsEach consumer's owning team
Compaction debtlsm.pending_compaction_bytes, write stallsIngest for affected nodesRate-limit backfills, add IOStorage/metrics platform
Resharding divergencemigration.row_diff_count > 0Users in migrated cohortRoll back per partitionDB platform
Unbounded cardinalitytsdb.active_series growth > 20%/dayWhole TSDB clusterDrop/aggregate offending labelMetrics 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.

OptionWhat WorksWhat BreaksWho Pays
Every team picks its storeFit-for-purpose, fast starts6 database engines, 3 sharding schemes, nobody expert in any; resharding learned per teamOn-call fatigue; security/compliance reviews ×6
One mandated database for all writesDeep expertise, single toolchainTelemetry crushes the OLTP fleet; records priced like telemetry or vice versaEveryone during the noisy-neighbor incident
Tiered platform: records / events / telemetry, each with a paved store, routing and reshardingTeams declare a tier; platform owns operations; costs visible per tierPlatform must support 3–4 engines well; exceptions needed for genuinely novel workloadsPlatform 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.

ScaleWrite Rate (records / telemetry)InfraPeople / On-callRough Monthly Total
Startup1K / 10K per sec1 Postgres primary + replica ($3K), small Kafka or managed queue ($1K)0.3 FTE; team on-call~$4K + ~$8K people
Growth15K / 500K per sec16-shard Postgres ($45K), Kafka 12 brokers ($15K), TSDB ($20K), object store ($3K)2 FTE DB platform; shared rotation~$85K + ~$50K people
Large200K / 10M per secDistributed SQL or 128 shards ($350K), Kafka multi-cluster ($120K), TSDB/columnar ($150K), object store ($30K)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#

Diagram: The 3-Year Evolution Path

One-Way Doors vs Two-Way Doors#

DecisionDoor TypeReversibility Cost
Batch size, linger, index setTwo-wayConfig or migration in hours
Number of physical shards (with logical partitions)Two-wayMove partitions; routine
Shard key choiceOne-wayRe-keying every table and every query; 2–4 quarters
Number of logical partitionsOne-way-ishChanging it is a full rehash; pick generously (1K–16K)
ID format without embedded routing infoOne-way-ishEvery lookup needs a directory; hard to retrofit
Exposing "write acknowledged = visible everywhere" in a public APIOne-wayClients depend on synchronous derivation forever
Telemetry vs record tier for a data typeTwo-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:

  1. 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.
  2. Every OLTP table carries the owning entity's partition key; no cross-partition foreign keys or global sequences.
  3. Derived data (search, analytics, caches, counters) MUST be produced from CDC/outbox, not from synchronous dual writes in request handlers.
  4. 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.
  5. 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#

SignalWhat 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 AsksWhat They Are TestingHow 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#

DimensionSenior (L5)Staff (L6)Principal (L7)
Framing"Shard it"Amplification, ack point, skew before partitioningClassifies org write volume into tiers with cost and durability contracts
PartitioningHash by IDLogical partitions, key chosen for spread + reads, hot-key planPlatform-owned routing and resharding as a routine operation
DurabilityImplicitAck point per path, loss budget signed offD0/D1/D2 tiers, mapped and enforced
FailureRetriesBackpressure, lag SLOs, compaction and skew alerts, reversible migrationsCells, rehearsed resharding, backfill governance
CostMore nodesHeadroom from fewer indexes and batching$/million writes per tier; retention as a cost lever

Strong Hire Signals

SignalWhat 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

SignalWhy It Misses the Bar
Replicas as write scalingFundamental misunderstanding of replication
Timestamp partition key for high-volume appendsGuarantees a hot head
No resharding storyShips 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#

NumberContext
10–50K/sSmall transactional writes on one tuned Postgres/MySQL primary (NVMe, group commit)
0.1–1 msNVMe fsync; the floor for a durable commit
+1–2 msAdding a synchronous in-region replica
~1,000 WCU/sDynamoDB per-partition write ceiling
10–50 MB/sPractical throughput per Kafka partition
50–100K/sWrites per node on Cassandra/Scylla-class LSM stores
< 100 MBHealthy 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–16KLogical 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
  1. Loading the index…