AlgoMaster Logo

Scaling Write Traffic

High Priority15 min readUpdated September 29, 2026
AI Mock Interview

Practice this topic in a realistic system design interview

Listen to this chapter
Unlock Audio

Premium Video

This video is available to premium subscribers only

Unlock Full Access

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.

Where This Pattern Shows Up

Let's start with the kinds of problems where the system has to handle a lot of writes.

ProblemWhat Gets Written
Design UberEvery active driver's phone sends its location every few seconds
Design WhatsAppEvery message sent is a new write, adding up to billions of messages per day
Design YouTubeEvery video view increases a view count, and every like updates a like count
Ad click trackingEvery 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.

1. Why Writes Are Harder to Scale Than Reads

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.

2. Start With the Database Itself

The first step is to make each write cheaper.

2.1 Keep Only the Indexes You Need

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.

2.2 Choose a Write-Friendly Database

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.

AspectB-Tree (PostgreSQL, MySQL)LSM Tree (Cassandra, ScyllaDB, RocksDB)
Write pathUpdates pages in place, scattered on diskAppends to memory and a log, flushes sorted files sequentially
Read pathFast point and range readsMay check several files to find a value
Background workLittleCompaction merges files over time
Good fitMixed workloads, transactionsHigh 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.

3. Batching Writes

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.

  • It adds some delay, because data waits in the buffer before it's written.
  • Buffered data can be lost. If the server crashes before the buffer is written, the data in that buffer is gone.

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.

4. Using a Queue to Absorb Writes

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.

4.1 Trade-offs

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.

5. Sharding

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.

5.1 Choosing the Shard Key

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 KeyWrite SpreadData Read TogetherVerdict
Driver ID (hashed)Even across millions of driversA driver's history is on one shardGood
Timestamp (ranges)All new writes hit the newest shardTime ranges are easy to scanPoor for writes
Random IDEvenRelated data is scatteredEven, 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.

5.2 Adding Shards Later

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.

6. Hot Shards and Hot Keys

Even with sharding, writes can still end up concentrated in one place.

6.1 A Bad Shard Key

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.

7. Reducing the Number of Writes

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.

7.1 Aggregate Before Writing

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.

7.2 Overwrite Instead of Adding Rows

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.

7.3 Ask Whether Every Write Needs Storing

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.

TechniqueExampleWhat You Give Up
AggregateWrite "views += 500" every few secondsA few seconds of counts if the aggregator crashes
OverwriteKeep only each driver's latest location in RedisThe full history, unless it's also stored elsewhere
SampleStore 1 in every 10 location pings for trip historyDetail between the stored points

8. Deciding Which Writes Can Be Lost

Many of the techniques in this lesson make writes faster by accepting some risk.

  • Batching in memory can lose data if a server crashes.
  • Aggregating view counts in memory can lose a few seconds of counts.
  • Acknowledging a write after putting it in a fast in-memory buffer, like Redis, and saving it to the database later (the write-behind pattern) can lose acknowledged writes if the buffer fails first.

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.

GroupExamplesConfirm to the User WhenTechniques That Fit
Must never be lostPayments, orders, chat messagesThe write is durable and replicatedReplicated database, queue with acks=all, idempotent retries
Small loss is acceptableView counts, analytics events, location pingsAs soon as it's acceptedIn-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.

9. Putting It All Together

Let's put these ideas together, using the Uber location example.

  1. Each driver's phone sends a location update every 4 seconds.
  2. The update reaches an API server, which adds it to a Kafka topic, partitioned by driver ID.
  3. One group of consumers reads from Kafka and updates each driver's latest location in Redis, overwriting the previous value.
  4. The matching service reads from Redis when it needs to find nearby drivers.
  5. Another group of consumers batches location updates and writes them to a Cassandra cluster, sharded by driver ID, for trip history and analytics.

Each piece handles one problem.

ComponentProblem It Solves
KafkaAbsorbs spikes
RedisKeeps only the current state, so 1.25 million updates per second overwrite 5 million entries instead of adding rows
Batching consumersReduce the number of database writes
Sharded CassandraSpreads 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.

10. Discussing High Write Traffic in an Interview

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.

QuestionWhy 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.

11. Key Takeaways

  1. Writes don't scale like reads. Caches and read replicas don't reduce the work of storing new data, and each write also updates logs and indexes.
  2. Make each write cheaper by keeping only the indexes you need and choosing a write-friendly database, such as an LSM-based store like Cassandra.
  3. Reduce the number of writes with batching, aggregation, and overwriting the latest value instead of adding rows.
  4. Use a queue to absorb spikes, and make consumers idempotent because messages can be delivered more than once.
  5. Use sharding to spread writes across many machines, with a shard key that spreads writes evenly and keeps related data together.
  6. Watch for hot spots. Avoid increasing values as range shard keys, and split counters for single very popular records.
  7. At each step, decide which writes can tolerate delay or loss, and which ones can't.

Quiz

Scaling Write Traffic Quiz

20 quizzes