Practice this topic in a realistic system design interview
Some systems generate new data continuously at a very high rate, such as driver locations, chat messages, user clicks, and application metrics. Every update becomes a write, and handling a large volume of concurrent writes is generally harder to scale than serving reads.
This lesson covers how to design a system that handles a very large number of writes, why writes are harder to scale than reads, and the main techniques used to handle them.
It also covers how to approach high write traffic in a system design interview, since this pattern comes up in many common interview problems.
Let's start with the kinds of problems where the system has to handle a lot of writes.
| Problem | What Gets Written |
|---|---|
| Design Uber | Every active driver's phone sends its location every few seconds |
| Design WhatsApp | Every message sent is a new write, adding up to billions of messages per day |
| Design YouTube | Every video view increases a view count, and every like updates a like count |
| Ad click tracking | Every click on an ad has to be recorded, because advertisers are billed based on those clicks |
| Metrics or logging (Datadog) | Thousands of servers send new data points every few seconds |
In all of these systems, new data keeps arriving all the time, and the database has to keep up with it.
With read traffic, we had a lot of options. We could cache results, add read replicas, and use a CDN. But most of these don't help with writes.
A cache can serve the same result to many readers, but every write carries new data that has to be stored. Read replicas don't help either, because every write still goes to the primary, and then every replica has to apply that same write too.
On top of that, each write does more work than it seems. The database has to write the change to a log on disk, update the table, update every index on that table, and sometimes lock rows so other writes don't conflict.
Let's put some numbers on this. Suppose Uber has 5 million active drivers, and each one sends a location update every 4 seconds. That's about 1.25 million writes per second.
No single database server can handle that. A well-tuned relational database handles somewhere in the tens of thousands of simple inserts per second, depending on hardware, indexes, and durability settings.
So we need ways to reduce the work each write does, and to spread writes across many machines.
The first step is to make each write cheaper.
One simple change is to remove indexes you don't need. Every index on a table has to be updated on every insert, so a table with five indexes does much more work per write than a table with one.
In a read-heavy system, extra indexes are worth it. In a write-heavy system, you should only keep the indexes your queries actually use.
The type of database also matters.
Traditional relational databases like PostgreSQL and MySQL store data in B-trees. When you insert a row, the database finds the right place in the tree and updates it, which can mean writing to many different places on disk.
Databases like Cassandra, ScyllaDB, and RocksDB use a different structure called an LSM tree (log-structured merge tree). They first write new data to memory and to a log file. When the in-memory data gets large enough, they write it to disk as one sorted file, in one sequential write.
Sequential writes are much faster than writes scattered across the disk, so LSM-based databases can usually handle much higher write throughput.
The trade-off is that reads can be slower, because the database may have to check several files to find a value. LSM databases reduce this cost with Bloom filters and a background process called compaction, which merges the sorted files over time, but compaction itself uses disk and CPU.
| Aspect | B-Tree (PostgreSQL, MySQL) | LSM Tree (Cassandra, ScyllaDB, RocksDB) |
|---|---|---|
| Write path | Updates pages in place, scattered on disk | Appends to memory and a log, flushes sorted files sequentially |
| Read path | Fast point and range reads | May check several files to find a value |
| Background work | Little | Compaction merges files over time |
| Good fit | Mixed workloads, transactions | High volumes of append-only data |
That's why databases like Cassandra are a common choice for write-heavy data like location history, messages, and metrics.
Another way to reduce work is to group many writes into one.
Each write to a database has a fixed cost, like a network round trip, a transaction, and a write to the log. If you insert 1,000 rows one at a time, you pay that cost 1,000 times. If you insert them in one batch, you pay it only once.
For example, a metrics service might collect data points in memory for one second, and then write all of them to the database in a single batch. This can increase write throughput by a large factor.
In practice, a batching buffer flushes on whichever comes first: the batch reaching a size limit, or a time limit passing. The size limit keeps batches efficient under heavy traffic, and the time limit makes sure data doesn't sit in the buffer for long when traffic is light.
Most databases support bulk inserts: a multi-row INSERT ... VALUES (...), (...), PostgreSQL's COPY, or a batch statement in Cassandra.
But batching has two trade-offs.
So batching works well for data like metrics or logs, where a short delay is fine and losing a small amount of data is acceptable. For something like a payment, you'd want each write to be stored safely before you tell the user it succeeded.
Write traffic is rarely steady. A flash sale, a big sports match, or a viral post can cause a sudden spike that's many times higher than normal traffic. If all of those writes go straight to the database, the database can get overloaded.
One way to handle this is to put a message queue, like Kafka, between the application and the database.
When a write comes in, the application adds it to the queue and responds to the client right away. Then a group of workers, called consumers, read from the queue and write to the database at a rate it can handle. During a spike, the queue grows, and the workers catch up once traffic goes down.
Kafka is well suited for this because it's built to handle very high write rates by appending messages to logs on disk. A Kafka topic is split into partitions, and messages with the same key, such as the same driver ID, always go to the same partition, so they stay in order.
This approach has a few trade-offs.
Delayed visibility. The data isn't in the database the moment the client gets a response, so a user might not see their change right away.
Duplicates. Most queues deliver each message at least once, so a worker might process the same message twice. For example, a worker writes an event to the database and then crashes before telling the queue it's done. The queue then delivers the same message to another worker.
To handle this, writes should be idempotent, which means applying the same write twice has the same result as applying it once. For example, each event can carry a unique ID, and the database can ignore an event if it has already stored that ID.
The backlog can keep growing. A queue handles a short spike, but if consumers are slower than producers for a long time, the backlog keeps growing. Watch consumer lag, the number of messages written but not yet processed. If it keeps rising, add consumers (up to the number of partitions) or larger batches. If even that isn't enough, the system needs backpressure: a bounded buffer that tells producers to slow down or retry later, instead of accepting work until memory runs out.
Batching and queues help a lot, but at some point, a single database still can't keep up. At that point, we need to split the data across multiple databases. This is called sharding.
Each database, called a shard, stores only part of the data, and handles only the writes for that part.
For example, in our Uber example, we could shard location data by driver ID. We hash each driver ID, and the hash decides which shard stores that driver's data.
If one shard can handle 50,000 writes per second, then 1.25 million writes per second needs at least 25 shards, and in practice more, to leave room for growth and spikes.
The most important decision here is the shard key.
A good shard key spreads writes evenly across all shards. It should also keep data that's read together on the same shard, so most queries only need to go to one shard.
Driver ID works well here, because writes are spread across millions of drivers, and when we need a driver's location history, it's all on one shard.
| Shard Key | Write Spread | Data Read Together | Verdict |
|---|---|---|---|
| Driver ID (hashed) | Even across millions of drivers | A driver's history is on one shard | Good |
| Timestamp (ranges) | All new writes hit the newest shard | Time ranges are easy to scan | Poor for writes |
| Random ID | Even | Related data is scattered | Even, but queries touch every shard |
A good shard key also keeps most writes inside one shard. A write that has to update two shards at once, like moving money between two accounts on different shards, needs extra coordination such as a saga or two-phase commit, so it's worth modeling data to avoid that for the common case.
With a simple hash(key) % N, changing the number of shards moves almost every key to a different shard. Many systems use consistent hashing instead: each shard is placed at a point on a ring, each key is stored on the next shard around the ring, and adding a shard moves only the keys in its new range, roughly 1/N of the data.
Even with sharding, writes can still end up concentrated in one place.
The first case is a bad shard key. Suppose you shard events by timestamp, and each shard holds a range of time. Every new event has the latest timestamp, so every new write goes to the same shard, while all the other shards sit idle.
That's why monotonically increasing values, like timestamps or auto-incrementing IDs, are usually a poor choice for range-based sharding. For time-series data, a composite key like (device_id, day) works better: the device ID spreads writes across shards, and the day keeps each device's data grouped for range scans.
The second case is a single very popular record.
Suppose a celebrity posts a video, and it gets 100,000 likes per second. Every like updates the same like count, which lives in one row on one shard. That row becomes a bottleneck, because updates to it can't run in parallel.
One common fix is to split the counter into several smaller counters. Instead of one like count for the video, you keep 10 counters, and each like updates one of them at random. When you need the total, you add the 10 counters together.
This spreads the writes across 10 rows, which can live on different shards, at the cost of a slightly more expensive read. Reads of a popular count are usually cached anyway, so the extra read cost is small.
So far, we've made writes cheaper and spread them out. But often, the best option is to write less data in the first place.
One way is to aggregate before writing. Take YouTube view counts.
Instead of updating the database every time someone watches a video, the system can count views in memory, or in a stream processor, for a few seconds. Then it writes one update, like "add 500 views", instead of 500 separate updates.
Metrics agents like StatsD use the same idea: they collect counters in memory and write one summary per interval. Time-series databases also keep recent data at full detail and combine older data into larger time buckets.
Another way is to overwrite instead of adding new rows.
In design Uber, the matching service only needs each driver's current location. So the latest location can be stored in an in-memory store like Redis, where each update simply replaces the previous one. Only a sampled or batched version of the location history needs to go to the permanent database.
In an interview, it's worth asking whether every single write really needs to be stored. Sometimes the answer is yes, like with payments or chat messages. But for things like view counts, metrics, and location pings, aggregating or sampling can reduce the write load a lot.
| Technique | Example | What You Give Up |
|---|---|---|
| Aggregate | Write "views += 500" every few seconds | A few seconds of counts if the aggregator crashes |
| Overwrite | Keep only each driver's latest location in Redis | The full history, unless it's also stored elsewhere |
| Sample | Store 1 in every 10 location pings for trip history | Detail between the stored points |
Many of the techniques in this lesson make writes faster by accepting some risk.
And even with a queue, you have to choose how safely each message is stored. For example, Kafka lets a producer wait until a message is copied to several brokers before it's confirmed (acks=all), which is safer but slower. Or it can wait for just one broker (acks=1), which is faster but can lose data if that broker fails.
So it helps to sort writes into two groups.
The first group is writes that must never be lost, like payments, orders, and chat messages. These should be stored durably, and ideally replicated, before you confirm them to the user.
The second group is writes where losing a small amount is acceptable, like view counts, analytics events, and location updates that will be replaced a few seconds later. For these, you can use faster techniques like in-memory batching and aggregation.
| Group | Examples | Confirm to the User When | Techniques That Fit |
|---|---|---|---|
| Must never be lost | Payments, orders, chat messages | The write is durable and replicated | Replicated database, queue with acks=all, idempotent retries |
| Small loss is acceptable | View counts, analytics events, location pings | As soon as it's accepted | In-memory batching, aggregation, overwriting, sampling |
In an interview, it's worth saying which group each type of write belongs to, since that decides which of these techniques you can use.
Let's put these ideas together, using the Uber location example.
Each piece handles one problem.
| Component | Problem It Solves |
|---|---|
| Kafka | Absorbs spikes |
| Redis | Keeps only the current state, so 1.25 million updates per second overwrite 5 million entries instead of adding rows |
| Batching consumers | Reduce the number of database writes |
| Sharded Cassandra | Spreads the remaining writes across many machines, using an LSM storage engine |
Because both consumer groups read from the same topic, each one can fall behind or restart without affecting the other, and because Kafka is partitioned by driver ID, each driver's updates are processed in order.
Here's a step-by-step approach you can use in an interview.
1. Clarify the requirements. Ask how many writes per second the system needs to handle, how spiky that traffic is, and whether any writes can be delayed or lost.
| Question | Why It Matters |
|---|---|
| How many writes per second at peak? | Tells you whether one database is enough |
| Is traffic steady or spiky? | Tells you whether you need a queue for bursts |
| Can any writes be delayed or lost? | Decides whether batching, aggregation, and async writes are allowed |
| Does every write need to be stored, or only the latest value? | Tells you whether overwriting or aggregation is possible |
| How will the data be read? | Affects the shard key and the database choice |
2. Estimate the peak write load, so you know whether one database is enough.
3. Look for ways to reduce the number of writes, using aggregation, overwriting, or batching.
4. Choose a database that fits the write pattern, such as an LSM-based database for high volumes of append-only data.
5. Add a queue if traffic is spiky or if the client doesn't need the write to be stored immediately.
6. Shard the data, explain your choice of shard key, and describe how you'd handle hot keys.
7. Explain which writes need strong durability, and which ones can trade some safety for speed.
In this order, each component you add is tied to a number or requirement from an earlier step.
20 quizzes