Pattern: Scaling Writes
Reads are easy to scale because copies are cheap: caches and replicas multiply read capacity. Writes are harder, because every write must reach the authoritative copy, and there is no cache for a change. When an interviewer says "now ten times the write traffic", you need a small set of techniques and the judgement to choose among them.
The unifying idea: spread writes across more machines, do fewer and larger writes, or avoid the write altogether. This chapter takes them in the order you would try them.
1. Understand the write load first
Questions that decide the approach:
- Rate and size. Writes per second, and bytes per write. A million tiny counter increments and a thousand large documents need different solutions.
- Does each write need to be durable and visible at once? A payment does. A page-view counter does not.
- Which key do they hit? Evenly spread keys scale by partitioning. A few hot keys do not.
- Do writes depend on reads? Read-modify-write needs coordination. Blind appends do not.
- Is ordering needed? Per key, globally, or not at all.
State these in the interview before choosing a technique. Then use the least complex tool that handles the numbers.
2. Reduce the write cost first
Before adding machines, make each write cheaper.
Remove unnecessary writes. Do not store what you can derive. Do not write a value that has not changed. Sample low-value events rather than recording all of them.
Trim indexes. Every index on a table is updated on every write. Remove ones that no query needs.
Batch. A thousand single-row inserts cost far more than one insert of a thousand rows, because each statement pays for a network round trip, parsing, locking and a log flush. Group writes, in the application or in a buffering layer, so that the database does fewer, larger operations.
Use the cheap storage structure. Append-only and log-structured stores turn random writes into sequential ones, which are much faster. If the write path is the bottleneck and updates are rare, a log-structured store, as in the wide-column chapter, may be the right engine.
Make the write path short. Do the minimum synchronously, such as persist the event, and defer the rest, such as updating derived views, to asynchronous consumers.
3. Buffer and batch with a queue or log
Put a durable queue or log between producers and the store. This does three things:
- It absorbs bursts. A spike fills the queue and is drained at the rate the store can sustain, instead of overwhelming it.
- It enables batching. Consumers read many messages and write them in bulk.
- It decouples the producer from the store's latency and availability.
The cost is that the write is not immediately visible, and you must handle failure: acknowledge to the producer only after the message is durably in the queue, make the consumer idempotent because a batch may be processed twice, and monitor queue depth, since a queue that grows forever is a delayed outage.
Where it fits: event ingestion, metrics, logs, analytics, notifications, and any write whose result nobody needs in the same millisecond. Where it does not: a write whose success the user must see confirmed with the data already readable. There, write synchronously, and use a queue only for the follow-up work.
Write-behind caching is a related idea: accept the write into a fast cache and flush to the database later. It makes writes fast and risks losing data if the cache fails before flushing, so use it only for data whose loss is tolerable, and with a durable log if you can.
4. Partition (shard) the data
When one node cannot sustain the write rate, split the data so that each node takes a share. This is the central scaling technique for writes.
Choose the shard key. It must:
- spread writes evenly,
- appear in the common queries, so reads also hit one shard,
- keep data that is written together in the same place.
Strategies.
- Hash partitioning spreads keys evenly and loses range queries.
- Range partitioning supports range queries, and sequential keys such as timestamps concentrate all new writes on one shard, a hot shard. Add a prefix, such as a bucket number or a source identifier, to spread them.
- Consistent hashing lets you add and remove shards while moving only a small fraction of keys.
Costs. Cross-shard queries and transactions become expensive or impossible, so design the data model so that the main operations stay within one shard. Rebalancing is operational work. Plan capacity, or use an engine that splits partitions automatically.
Alternatives to sharding a relational store. Use a partitioned engine built for the workload, such as a wide-column or distributed key-value database, where partitioning is automatic and writes scale linearly with nodes, if your access patterns fit its model.
5. Hot keys and hot partitions
Partitioning spreads keys, not traffic to one key. If one key gets a huge share of writes, such as a celebrity's follower counter, the global like count on a viral post, or a flash-sale inventory row, the shard that owns it is the bottleneck, however many shards exist.
Detect it. Track write rates per key or per partition, and alert on skew.
Techniques.
- Split the key (salting). Write to one of sub-keys chosen at random or by a hash of the writer, such as
likes:post42#0to#7, and sum them on read or merge them periodically. Writes scale by , and reads cost lookups, or a periodically refreshed total. - Aggregate before writing. Buffer increments in memory in each application server and flush a combined update every second. A thousand increments become one write.
- Approximate. If an exact count is not required, use an approximate counter or sample, and show "about 1.2 million".
- Move contended updates to a queue that processes them in order, one worker per key, so the writes to the key are serial and never conflict. (See the contention pattern.)
- Avoid the shared counter altogether. Derive the count from the set of events when needed, or keep a per-user record and aggregate asynchronously.
State the trade-off: salting and aggregation give up exactness and immediacy for throughput.
6. Reduce coordination
Writes slow down when they must agree with each other.
- Prefer commutative, append-only operations. Adding an event to a log needs no knowledge of the other events. "Set a value to the sum of what I read" does.
- Use idempotent writes, so retries and duplicates are safe, which allows at-least-once pipelines and simple recovery.
- Choose weaker consistency where it is acceptable. Asynchronous replication and quorum-light writes are faster than synchronous replication to every node.
- Keep transactions small and local. A transaction spanning shards needs a coordination protocol that is slow and fragile.
- Single-writer per key. If only one process writes a given key, there is no contention to resolve.
7. Handling the failure cases
- The queue or log fills up. Apply backpressure to producers, shed low-priority writes, and scale consumers. Never let it grow without bound.
- A shard is overloaded or down. Replicate each shard, and decide how a write to an unavailable shard behaves: fail, buffer and retry, or write elsewhere for later reconciliation.
- Duplicates. Make writes idempotent with a natural or generated identifier.
- Partial batch failure. Make the batch retryable as a unit, or track per-record status.
- Ordering broken by retries or partitions. Order per key by partitioning on the key, and attach sequence numbers if consumers need to detect gaps.
- Lost writes in a failover. With asynchronous replication, the last writes may be lost. If that is unacceptable, use synchronous or quorum replication for that data, and accept the latency.
8. A worked example
Problem. A game records player events. Peak load is 200,000 events per second, each about 200 bytes, with a per-player running score that updates on many events, and a global event count shown on the home page.
Reasoning.
- 200,000 writes per second at 200 bytes is 40 MB per second, which is far beyond what a single relational primary should absorb as individual transactions.
- Ingest through a partitioned log. Producers write events to a topic keyed by player. The write is acknowledged once it is durably in the log. This absorbs spikes.
- Consumers batch events into the event store, a log-structured, partitioned store keyed by player and time bucket, as bulk writes. This turns 200,000 small writes into a few thousand larger ones.
- Per-player score. All events for a player land in one partition and are processed in order by one consumer, so the score is updated by a single writer with no contention. Store the result in a key-value store.
- Global count. One counter would be a hot key. Maintain 16 salted counters updated from aggregated in-memory counts every second, and sum them for the home page, which tolerates being a second or two stale.
- Idempotency. Each event has an identifier, so a re-delivered batch overwrites the same rows.
What I would say about the trade-off. "Scores and counts are eventually consistent by a second or two. I trade immediacy for the ability to absorb the peak, and the event log lets me rebuild any derived value if a consumer has a bug."
9. Interview questions and model answers
Q: Writes are the bottleneck. What do you try first? Make each write cheaper: remove unneeded indexes and writes, batch, and shorten the synchronous path. Then buffer through a queue to absorb bursts, and then partition by a key chosen from the main access pattern.
Q: How do you handle a hot key? Split it into salted sub-keys and sum on read, aggregate increments in memory before writing, serialise updates through a queue with one worker per key, or use an approximate count, depending on how exact it must be.
Q: What does a queue in front of the database cost you? The write is not immediately visible, I must make consumers idempotent and handle backpressure, and I need to monitor queue depth. In return I absorb bursts and can batch.
Q: How do you choose a shard key? So that writes spread evenly, the common queries include it, and data written together stays together. I avoid monotonically increasing keys on a range-partitioned store because they create a hot shard.
Q: How do you avoid cross-shard transactions? By choosing the shard key so that the main operations stay within one shard, and using sagas or compensating actions for the rare flows that cross shards.
10. Common mistakes
- Sharding before reducing the cost of each write and trying batching.
- A monotonically increasing shard key on range partitioning, causing a hot shard.
- A single global counter updated by every request.
- A queue with no backpressure or monitoring.
- Non-idempotent consumers on an at-least-once pipeline.
- Forgetting that salted counters need a read-time merge.
- Synchronous replication to every node when the data does not need it.