Skip to content

The finished design, decision by decision

How to design Discord's Message Storage

Not the one correct diagram, but a design you can defend under these constraints: the finished architecture, then every stage's question with the reasoning that answers it, the tradeoffs it accepts, and where another engineer could land differently.

Storing trillions of chat messages

The short answer

5 parts, each with one job. The map below shows how requests and data move between them; the stages after it explain why each part is there.

Members
Send messages, open channels, scroll history.
API servers
Validate, assign Snowflake IDs, check permissions, and call the message service.
Gateway
Holds members' WebSocket connections and pushes new-message events.
Message data service
The only path to the database. Routes each channel to one instance by consistent hashing and coalesces identical concurrent reads.
Message cluster
Wide-column, log-structured store. Partition key (channel_id, bucket); rows clustered by message_id, newest first; replication factor 3.
12345CLIENTMembersSERVICEAPI serversSERVICEGatewaySERVICEMessagedata serviceDATABASEMessage cluster

Select a component to see what it is responsible for and which state it owns.

  1. 1Members → API servers: Send, load history, jump
  2. 2API servers → Gateway: New message event
  3. 3Gateway → Members: Push to online members
  4. 4API servers → Message data service: Query by channel (hash-routed)
  5. 5Message data service → Message cluster: Read and write one partition
  • Request / response
  • Asynchronous
  • Server push

Why does this design work?

Every common read is shaped to touch one small, contiguous partition: messages are keyed by (channel, ten-day bucket) and ordered by a time-sortable Snowflake ID, so the latest page, a history page and a jump to any message each map to a known partition. A log-structured, replicated store makes writes cheap appends and grows by adding nodes.

Its weaknesses are handled explicitly: tombstones are kept short-lived, nulls are never written, and empty buckets are skipped; half-rows from edit/delete races are recognised and removed. A data service in front of the store routes each channel to one instance and coalesces identical reads, so a hot channel costs one query. The same layer made it possible to dual-write and move the whole dataset to a new database while it was live.

Invariants, and where they are enforced

  • Loading a channel reads one bounded partition, however large or old the channel is.

    The partition key includes a fixed time bucket derived from the message ID, so partitions stay small; empty buckets are tracked and skipped. Enforced by Message cluster, Message data service.

  • A burst of readers on one channel becomes one database query, not thousands.

    Consistent hashing sends every request for a channel to the same instance, which merges identical in-flight reads. Enforced by Message data service.

  • An acknowledged message survives the loss of a database node.

    Writes go to three replicas and are acknowledged at quorum. Enforced by Message cluster.

What does it rely on?

  • Almost every query is by channel and time.
  • Snowflake IDs are generated with roughly synchronized clocks, so they sort by time.
  • Repair runs regularly, so short tombstone lifetimes are safe.
  • Requests for a channel can be routed to one data-service instance.

What tradeoffs does it make?

ChoiceGainsCosts
Wide-column LSM storeCheap writes, scale-out, built-in replication.Narrow queries, tombstones, compaction and repair to operate.
Time-bucketed partitionsBounded partitions whatever the channel's size.Quiet channels need several bucket reads per page.
Data service with coalescingHot channels cost one query; one place for limits and migrations.Another service to run; routing must stay consistent.

What are the reasonable alternatives?

Sharded Postgres by channel
Better when growth is slower, relational queries matter, or the team knows Postgres deeply.
Managed key-value store (DynamoDB, Bigtable)
Better when you want the same data model without operating the cluster yourself.
Tiered storage: hot recent buckets in the database, old buckets in object storage
Better when most history is never read again and storage cost dominates.

When does it stop working?

  • Queries need to filter by something other than channel and time (search, mentions), which calls for separate indexes fed from the message stream.
  • A single channel's write rate exceeds what one partition's replicas can absorb.
  • Clock skew between ID generators is large enough to misorder messages noticeably.

Every stage, decided and explained

Spoilers, for the whole investigation: each stage's question and its answer, the reasoning behind it, and the tradeoffs it accepts. If you have not worked through the stages yet, you may want to do that first.

Work through the stages

Stage 1 of 9 · Model

What the numbers say

Use 86,400 seconds in a day. Messages are about 1 KB.

What you need to know first

120 million messages a day. About how many writes a second on average?

About 1,400 per second.

120,000,000 ÷ 86,400 ≈ 1,400 a second. Even several times that at peak is within one good database's reach. The write rate isn't the problem.

At about 1 KB per message, how many terabytes a year (before replication)?

About 44 TB.

120 GB a day × 365 ≈ 44 TB a year, and growing several-fold a year. What outgrows one machine is the data.

When data and indexes fit in memory, a random read is a memory lookup. When they don't, it becomes a disk seek, and latency becomes unpredictable. The usual fix is to store what one query needs physically together, so it takes one seek and a short sequential read instead of fifty random ones.

A Snowflake ID is 64 bits: a millisecond timestamp in the high bits, then a worker number and a per-worker sequence. Any server can generate one without coordinating, and sorting IDs sorts messages by creation time. See Generating unique identifiers.

So "the latest 50 messages in a channel" is "the 50 largest IDs in that channel".

Do time-ordered IDs require one central counter?

No: putting the timestamp in the high bits makes independently generated IDs sort by time.

Each worker's number and sequence keep IDs unique; the timestamp makes them sortable. No coordination needed.

What the stage asks

Which statements follow from the scenario?

  1. Holds

    120 million messages a day is roughly 1,400 writes a second on average.

    120,000,000 / 86,400 ≈ 1,390. Peaks are several times that, which a single good database could still absorb. The write rate is not the problem.

  2. Holds

    What outgrows one machine is the data, not the request rate.

    At ~1 KB a message that is ~120 GB a day, over 40 TB a year before replication, and growing. Once data and indexes stop fitting in memory, random reads go to disk and latency becomes unpredictable, which is exactly the symptom in the scenario.

  3. Depends

    Because reads are about half the traffic, a cache of recent messages will take most read load off the database.

    It helps for the latest page of busy channels. But reads are scattered across millions of small channels, history pages, jumps to old messages and mentions, so much of the read traffic is a long tail a cache will mostly miss. The store itself must serve random reads well.

  4. Holds

    The dominant query is 'latest messages in channel X', so a channel's messages should be stored together, in time order.

    If a page of messages is physically contiguous, loading it is one seek and a short sequential read rather than fifty random lookups.

  5. Fails

    To sort messages by time, IDs must come from one central sequence.

    Snowflake IDs put a millisecond timestamp in the high bits, then a worker number and a per-worker sequence. Any server generates them without coordination, and they sort by creation time to the millisecond. See Generating unique identifiers.

The reasoning

  1. The write rate is modest; what outgrows one machine is the ever-growing data.
  2. Store a channel's messages together in time order so the latest page is one contiguous read.
  3. Snowflake IDs sort by time without central coordination.

Chat storage is a data-shape problem. The write rate is modest; the dataset is huge, ever-growing and read at random. The design has to make the common read (a channel's latest page) touch a small, contiguous piece of data, whatever the total size.

Time-sortable IDs are what make that possible: if the ID is the time, "the latest 50" is just "the 50 largest IDs in this channel".

Stage 2 of 9 · Decide

Choose the store

The requirements ask for growth by adding nodes, survival of node loss, and a team too small for constant manual operations.

What you need to know first

A log-structured (LSM) store turns every write into an append: new data goes to memory, then is flushed to immutable sorted files on disk, which are periodically merged (compaction). Writes are cheap; reads may consult several files. See Log-structured storage (LSM trees).

A wide-column store like Cassandra or ScyllaDB has two keys per table:

  • The partition key decides which nodes hold a row, and which rows are stored together.
  • The clustering key orders rows within a partition.

Capacity grows by adding nodes; each partition is replicated (typically 3 copies) automatically. See Partitioning.

Why not one document per channel holding an array of its messages?

Documents grow without bound, every send rewrites a growing document, and size limits cap history.

The busiest channels become the slowest writes, and a 16 MB document limit eventually caps a channel's history.

What do you give up with a wide-column LSM store compared with Postgres?

Flexible queries (you query by key, not arbitrary filters or joins), strong consistency by default, and cheap deletes: a delete is a tombstone that reads must skip until compaction removes it. You also take on compaction and repair as ongoing operational work.

What the stage asks

Which store should hold the messages?

  1. Defensible

    Postgres, with an index on (channel_id, message_id), partitioned and sharded by hand as it grows

    It will work for a long time, and SQL is flexible. But at this growth rate you will be sharding by hand repeatedly, moving data between shards and managing failover yourself, which is what the 'add nodes, no manual resharding' requirement rules out for this team.

  2. Sound

    A wide-column, log-structured store (Cassandra-style) partitioned by channel and clustered by message ID

    Capacity grows by adding nodes; replication and failover are built in; writes are appends, which suits a write-heavy log of messages; and a partition's rows are stored together in clustering order, so a page of a channel is one contiguous read. The costs are a narrow query model, eventual consistency and real operational work (compaction, repair). This is what Discord chose in 2017.

  3. Flawed

    A document store with one document per channel holding an array of its messages

    Documents grow without bound, every send rewrites or locks a growing document, and document size limits (16 MB in MongoDB) cap a channel's history. It turns the busiest channels into the slowest writes.

  4. Flawed

    Redis as the primary store, with periodic snapshots to disk

    Tens of terabytes a year in RAM is very expensive, and snapshot persistence loses every message acknowledged since the last snapshot when a node dies, which violates the durability requirement.

What a strong answer covers

  • A channel's messages must be stored together in time order so the latest page is one contiguous read.
  • Growth by adding nodes with built-in replication matches the requirement and the team size.
  • Names the costs: limited queries, eventual consistency, compaction and repair to operate.
  • Notes that log-structured writes suit a write-heavy, append-mostly workload.Supporting

The reasoning

  1. Append-mostly data read by contiguous key ranges suits a wide-column LSM store.
  2. Growth by adding nodes, with built-in replication, matches a small team and fast growth.
  3. LSM costs: reads touch several files, deletes are tombstones, compaction competes with traffic.

The choice follows from three facts: the data is append-mostly, the main query reads a contiguous slice of one channel, and the team wants growth to be "add a node". A wide-column store built on Log-structured storage (LSM trees) fits all three.

It is not free. Reads cost more than writes in an LSM engine, deletes become tombstones, and compaction competes with traffic. Those costs are where the later failures in this investigation come from.

Where another engineer could land differently

A team with deep Postgres expertise and a slower growth curve could reasonably shard Postgres by channel and accept the manual work; several large chat products do.

Stage 3 of 9 · Decide

Choose the partition key

In this store, the partition key decides which nodes hold a row and which rows are stored together. The clustering key orders rows inside a partition. Some channels will receive millions of messages over their lifetime; most receive a few hundred.

What you need to know first

Partitions must stay bounded. A partition that grows forever (many gigabytes) makes compaction, repair and replacing a node slow and memory-hungry, and reads of it get slower.

The natural key here, the channel, is unbounded: a busy channel gets messages for years.

Bucketing bounds it: add a time window to the partition key, say 10 days. Partition = (channel, bucket), where the bucket is computed from the message's timestamp. Each partition holds at most 10 days of one channel.

Because the bucket comes from the Snowflake ID's timestamp, any message's partition can be computed from its ID alone.

A very busy channel gets 10,000 messages a day at 1 KB each. About how many megabytes in one 10-day bucket?

About 100 MB.

10,000 × 10 days × 1 KB = 100 MB, which is Discord's target ceiling. Without buckets, after 5 years the same channel's partition would be about 18 GB.

Why not partition by message_id so writes spread perfectly evenly?

The latest 50 messages of a channel would be on up to 50 different partitions: every page becomes a scatter-gather.

Even spreading helps writes, but the dominant read needs a channel's messages together.

What the stage asks

What should the primary key of the messages table be?

  1. Flawed

    Partition by channel_id; cluster by message_id descending

    Every page of a channel is one partition, which is good, but a busy channel's partition grows forever. Partitions of many gigabytes make compaction, repair and node replacement slow and memory-hungry, and reads of a huge partition get slower over the years. It fails the 'whatever the channel's age' requirement.

  2. Sound

    Partition by (channel_id, bucket), where bucket is a fixed time window (about 10 days) computed from the message ID; cluster by message_id descending

    A partition holds at most ten days of one channel, so even the busiest stays bounded (Discord aimed for under 100 MB). The latest page is a read of the current bucket; scrolling back walks to earlier buckets; jumping to a message computes its bucket from its ID. This is the schema Discord described.

  3. Flawed

    Partition by message_id, so writes spread perfectly evenly

    Writes spread well, but 'the latest 50 messages in a channel' now touches up to 50 partitions on different nodes. The common read becomes a scatter-gather.

  4. Flawed

    Partition by server (guild) id; cluster by channel_id, message_id

    The largest communities become enormous partitions holding every channel's history, and their traffic all lands on the same three replicas.

What a strong answer covers

  • A channel's messages for a time range are in one partition, ordered by message ID.
  • Partitions are bounded in size regardless of how busy or old a channel is.
  • The bucket is computed from the Snowflake ID's timestamp, so any message's partition is known from its ID.
  • Paging back means reading earlier buckets, possibly several if a channel is quiet.Supporting

The reasoning

  1. Keep a page's rows together and every partition bounded.
  2. Bucket an unbounded key by time: partition by (channel, bucket), cluster by message ID.
  3. Quiet channels may need to walk back through several buckets to fill a page.

PRIMARY KEY ((channel_id, bucket), message_id) does two jobs: it keeps the rows a page needs together, and it keeps every partition bounded. Bucketing by time is the standard way to stop a partition growing forever when the natural key (the channel) is unbounded.

The cost shows up for quiet channels: their latest 50 messages may be spread over many buckets, and a read walks backwards through them. Hold that thought for the next stage.

Tradeoffs

ChoiceGainsCosts
Bucket of ~10 daysBounded partitions; the latest page is usually one read.Quiet channels need several bucket reads to fill a page.
Cluster by message_id descendingThe newest messages are first on disk.Queries must be by channel and bucket; no ad-hoc filtering.

Stage 4 of 9 · Break it

The channel that froze the cluster

Here is what the message service and the store logged when one member opened the channel. Find the lines that explain the stall.

What you need to know first

In an LSM store, a delete is a write: a tombstone marking the row as deleted. Older copies of the row may sit in files that haven't been compacted yet, so the tombstone must be kept, and read past, until compaction removes both.

Tombstones are kept for a grace period (gc_grace_seconds) so replicas that missed the delete can learn of it through repair.

A channel had 2 million messages; a bot deleted all but one. A read asks for the latest 50 messages. How much work is it?

Huge. The read has to scan past every tombstone in each bucket before concluding the bucket has nothing live, then move to the previous bucket and do it again. Reading "nothing" costs millions of steps, on every replica that serves the read.

A writer inserts every column, writing explicit NULL for the 12 columns a message doesn't use. What does that do in this store?

Creates a tombstone for each null column: a dozen useless tombstones per message.

Writing NULL means 'delete this cell'. Write only the columns a message actually has.

Three ways to bound the damage: track which buckets are empty so reads skip them; shorten the tombstone grace period (safe if repair runs more often than the period); and stop writing needless nulls.

What the stage asks

Select the lines that are part of the problem.

LogOpening #announcements (channel 81723…)
  1. 1schema: messages has 16 columns; the average message sets 4 of them
  2. 2writer: INSERT INTO messages (…all 16 columns…) VALUES (…, null, null, null, …)

    Writing explicit nulls creates a tombstone for every null column. Most messages carried a dozen tombstones they never needed. Write only the columns that have values.

  3. 3config: gc_grace_seconds = 864000 # tombstones kept 10 days before compaction may drop them

    Ten days of tombstones from millions of deletes are still on disk. Discord shortened this to two days, which is safe because repair runs nightly and every replica sees deletes well within that window.

  4. 410:03:15 GET /channels/81723/messages?limit=50
  5. 510:03:15 read (81723, bucket 1871): live rows 0, tombstones scanned 211,903

    Reading a partition of deleted rows means scanning every tombstone in it before concluding it is empty. The work is proportional to what was deleted, not to what is returned.

  6. 610:03:16 read (81723, bucket 1870): live rows 0, tombstones scanned 198,277 … (continues back through 140 buckets)

    The reader walks backwards bucket by bucket looking for 50 live messages, through every empty bucket. Tracking which buckets are empty lets reads skip them.

  7. 710:03:24 node-17: JVM GC pause 9,870 ms
  8. 810:03:24 node-09, node-23: JVM GC pause 9,410 ms, 10,120 ms
  9. 910:03:25 cluster p99 read latency 4,200 ms

What the fix has to do

  • Deletes are tombstones that reads must scan until compaction removes them.
  • The read walks through many empty buckets; tracking and skipping empty buckets bounds it.
  • Shortening tombstone lifetime (safe given regular repair) reduces how many are kept.
  • Writing nulls creates needless tombstones; write only present columns.Supporting
  • All replicas of the partition do the same work, so the stall spreads to three nodes.Supporting

The reasoning

  1. In an LSM store a delete is a tombstone, and reads must scan past tombstones until compaction.
  2. Bound the work: skip empty buckets and shorten tombstone lifetime when repair runs regularly.
  3. Write only the columns you have; explicit nulls create tombstones.

This is the incident from Discord's 2017 post. In a log-structured store, a delete is a write: a tombstone that must be read past until compaction removes it. A channel with millions of deletions becomes a channel where reading "nothing" costs millions of steps, and those steps allocate enough memory to trigger stop-the-world garbage collection on every replica holding the partition.

Discord's fixes were all about bounding that work: skip empty buckets, keep tombstones two days instead of ten (nightly repair makes that safe), and stop writing nulls. See Log-structured storage (LSM trees).

Stage 5 of 9 · Break it

The message with no author

In this store, UPDATE and INSERT are both upserts: they write the given columns with a timestamp. Conflicts are resolved per column by last write wins.

What you need to know first

In this store, writes are blind upserts: an UPDATE doesn't check whether the row exists; it just writes the given columns with a timestamp. Conflicts are resolved per column by last write wins (LWW): each cell keeps its newest value. See Conflict resolution and convergence.

A row is a collection of cells, not a single value.

A delete at time T1 writes a tombstone for the row. An edit at T2 > T1 writes the body and edited_at columns. What does a reader see?

A half-row: body and edited_at (newer than the tombstone, so they win), and every other column deleted. A message with no author, channel or timestamp, which no user ever wrote.

The race is rare. What's the cheaper fix?

Treat a row missing a required column (like author_id) as deleted, and clean it up on read

The invalid state is easy to recognise and rare, so repairing it costs almost nothing.

What the stage asks

Which statements hold?

  1. Fails

    An UPDATE of a row that was just deleted fails, because the row no longer exists.

    Writes are blind upserts; there is no read first. The edit writes its columns whether or not a row exists.

  2. Fails

    Because the delete happened, last-write-wins guarantees the message stays deleted.

    Last write wins per column. The delete's tombstone covers the row at time T1; the edit writes body and edited_at at T2 > T1, and those newer cells win. The result is a half-row: the edited columns, with everything else deleted.

  3. Holds

    A row missing a required column such as author_id can be treated as a deleted message and cleaned up on read.

    This was Discord's fix: a message without an author cannot be valid, so readers treat it as deleted and remove it.

  4. Holds

    Making edits conditional (UPDATE … IF EXISTS) would prevent the half-row, at the cost of a consensus round trip on every edit.

    A lightweight transaction checks existence atomically, but it costs several round trips. For a rare race, detecting and repairing invalid rows is the cheaper trade.

The reasoning

  1. Per-column last-write-wins means each cell keeps its newest value, not that the last operation wins.
  2. A racing edit and delete can leave a row nobody wrote.
  3. When a race is rare and its result recognisable, repairing on read beats preventing with consensus.

Eventual consistency with per-column last-write-wins does not mean "the last operation wins". It means each cell keeps its newest value, and a row is just a collection of cells. An edit racing a delete leaves a row that no single user ever wrote.

You can prevent the race with a conditional write (Concurrency control) or accept it and repair what it produces. When the race is rare and the invalid state is easy to recognise, repair is usually cheaper. See Conflict resolution and convergence.

Stage 6 of 9 · Change it

Everyone opens the same channel

The cluster has plenty of total capacity. One partition is getting far more reads than three nodes can serve.

What you need to know first

A partition lives on its replicas, typically three nodes. Adding nodes adds capacity for other partitions; it can't split one partition's load. A burst of reads for one channel lands on the same three nodes however big the cluster is.

But hundreds of thousands of "latest page of channel X" requests are the same question. If they meet in one place, they can be answered by one query. That's request coalescing: the first request runs the query; identical requests that arrive while it's running wait for its result. See Request coalescing.

Requests only meet if they're routed to the same place, which is what consistent hashing on the channel ID does for a fleet of data-service instances.

The same arithmetic applies to a burst of identical reads: compare no protection, coalescing per instance, and one query fleet-wide.

A hot key expires. Rebuilding it takes one expensive query. Until that query finishes, every request for the key misses. The database can run about 200 of these queries at once before it slows down.
5,000/s400 ms

Queries reaching the database while the key is rebuilt (log scale). The dashed line is what it can run at once.

  • No protection2,000
  • Coalesce per server20
  • One rebuild, fleet-wide1
  • Serve stale, refresh behind1
No protection:
Every request that misses queries the database. 2,000 requests wait about 400 ms, longer, because the database is past capacity and every query slows down.
Coalesce per server:
Each of 20 servers lets one request rebuild; the rest on that server wait for it. 2,000 requests wait about 400 ms.
One rebuild, fleet-wide:
A short lock in the cache lets one request in the whole fleet rebuild. 2,000 requests wait about 400 ms.
Serve stale, refresh behind:
Keep serving the old value while one background request rebuilds it. Nobody waits; a few requests see a value that is a moment old.

Why route each channel's requests to one data-service instance?

So all identical requests for a channel reach the same instance, where they can be coalesced

Spread randomly over 50 instances, each would issue its own query. Routed by channel, one instance sees them all.

What the stage asks

What do you change?

  1. Flawed

    Add more nodes to the cluster

    The partition still lives on the same three replicas. More nodes add capacity for other keys; they do not split one key's load. See Partitioning.

  2. Sound

    Route every request for a channel to one data-service instance (consistent hashing on channel_id) and coalesce identical in-flight reads into one query

    Hundreds of thousands of identical 'latest page of channel X' requests arrive at the same instance, which issues one query and hands the result to every waiter. The database sees a handful of reads instead of a flood. Discord built exactly this layer in Rust in front of its database.

  3. Defensible

    Cache each channel's latest page in Redis with a five-second TTL

    It absorbs this burst. But every new message, edit and delete in a busy channel invalidates the page, so you are maintaining a second copy under constant churn, and when the entry expires the burst still stampedes the database unless the refill is coalesced.

  4. Defensible

    Read at consistency ONE from any replica instead of QUORUM

    It spreads reads over all three replicas instead of two per read, a modest gain. The burst is still far more than three nodes can serve, and you give up read-your-writes for the sender.

What a strong answer covers

  • A single partition's load cannot be spread by adding nodes; it stays on its replicas.
  • The burst consists of identical reads, so it can be served by one query.
  • Routing by channel is what lets one instance see, and merge, all the identical requests.
  • Mentions protecting unrelated channels on the same nodes.Supporting

The reasoning

  1. One partition's load stays on its replicas; adding nodes doesn't split it.
  2. A burst of identical reads is one question: coalesce it into one query.
  3. Route by channel (consistent hashing) so identical requests meet where they can be coalesced.

A hot key is a key-level problem, so it needs a key-level fix. The burst is not a thousand different questions; it is one question asked a thousand times. Request coalescing answers it once.

Coalescing only works if the duplicates meet, which is why it is paired with Consistent hashing on the channel ID: every instance sees all the traffic for the channels it owns. Discord's 2023 post credits this data-service layer, together with the database change in the next stage, for keeping the cluster calm through traffic spikes like the 2022 World Cup final.

Stage 7 of 9 · Change it

Move trillions of messages, live

Messages keep arriving throughout. The old cluster must stay correct until the very end.

What you need to know first

A live migration follows the same shape every time. See Online data migrations:

  1. Capture new writes in both places, so the new store never falls further behind.
  2. Backfill the past, which has now stopped changing.
  3. Verify by comparing reads from both.
  4. Move reads to the new store, keeping a way back.
  5. Stop writing to the old store.

Why start dual writes before copying the old data?

So everything after a known point is already in both stores; the copy only has to handle a past that no longer changes.

If you copy first, messages written during the copy are missing from the new store, and you have to chase a moving target.

Copying trillions of rows takes days. Why checkpoint each token range?

So a failure partway costs one range, not the whole copy. Without checkpoints, a crash on day 3 means starting over.

What the stage asks

Put the migration steps in a safe order.

In this order

  1. 1Stand up the new cluster and make the data service write every new message to both clusters — The old cluster remains the source of truth.
  2. 2Record the point from which every message exists in both clusters
  3. 3Copy everything older than that point, token range by token range, checkpointing each range — So a crashed migrator resumes rather than restarts.
  4. 4Send a sample of reads to both clusters and compare the results
  5. 5Switch reads to the new cluster, keeping dual writes so you can switch back
  6. 6Stop writing to the old cluster and decommission it

Dual writes come first so that nothing written during the copy is lost; the copy then only has to cover a fixed, unchanging past. Verification before switching reads, and dual writes kept on after it, keep every step reversible until the last one. See Online data migrations.

Discord's numbers: their first migrator (on Spark) estimated three months. They rewrote it in Rust in an afternoon, reached 3.2 million messages a second, and copied everything in nine days. The last few token ranges stalled on enormous runs of tombstones (the previous stage's problem, again), which one compaction cleared. The new cluster runs on 72 nodes instead of 177, with p99 history reads of 15 ms instead of 40 to 125 ms.

The reasoning

  1. Capture new writes first, then backfill a past that has stopped changing.
  2. Verify by comparing reads, then move reads before writes so there's always a way back.
  3. Checkpoint the copy per range and throttle it so live traffic isn't starved.

Every safe migration is the same shape: capture new writes first, backfill a past that has stopped changing, verify, then move reads before writes so there is always a way back.

The interesting engineering is in the details: checkpointing per token range so failures cost minutes rather than days, throttling the copy so it does not starve live traffic, and expecting the old store's worst data (here, tombstones) to be the last and slowest thing you move.

Stage 8 of 9 · Break it

Write the coalescer

Implement the read path in the data service so that identical concurrent reads share one database query. Permissions have already been checked by the API servers. Think about what happens when the shared query fails or hangs.

What you need to know first

A coalescer is a map from request key to in-flight promise:

  1. If the key is in the map, await the existing promise.
  2. Otherwise, start the query, store its promise, and remove the entry when it settles.

The key must include every parameter that changes the result: channel, cursor and limit. Two requests that differ in any of them aren't the same question.

The shared query fails. What must happen to its map entry?

Remove it, so the next request tries again instead of receiving the cached failure

Keeping a failed promise in the map would hand the same error to every later request.

The shared query hangs for 60 seconds. What happens to the hundreds of thousands of waiters, and what prevents it?

They all wait 60 seconds: coalescing turned one slow query into a slow page for everyone. A timeout on the shared query bounds it; when it fires, the entry is removed and the next request starts fresh.

What the stage asks

Implement getMessages with request coalescing.

Reference implementation

const inflight = new Map<string, Promise<Message[]>>();

function withTimeout<T>(p: Promise<T>, ms: number): Promise<T> {
  return Promise.race([
    p,
    new Promise<T>((_, reject) => setTimeout(() => reject(new Error("query timed out")), ms)),
  ]);
}

export function getMessages(channelId: string, before: string | null, limit: number): Promise<Message[]> {
  const key = `${channelId}:${before ?? "latest"}:${limit}`;
  const existing = inflight.get(key);
  if (existing) return existing;

  const query = withTimeout(queryPage(channelId, before, limit), 2_000).finally(() => {
    inflight.delete(key);
  });
  inflight.set(key, query);
  return query;
}
  • finally removes the entry whether the query succeeds or fails. Without it, one failed query would be returned to every future caller.
  • The map holds a request only while it is in flight. Results are not kept, so there is nothing to invalidate when a new message arrives.
  • The key must contain every parameter. Coalescing "latest 50" with "latest 100" would hand some callers the wrong page.
  • Per-user checks happen before this function. Never coalesce a response that depends on who is asking.
  • This only collapses requests that reach the same process, which is why the router hashes channel IDs to instances. See Request coalescing.

What a strong answer covers

  • Keeps a map of in-flight queries; later identical calls await the existing promise.
  • The key includes every parameter that changes the result (channel, cursor, limit).
  • Removes the entry when the query settles, on failure as well as success, so errors are not cached.
  • Bounds the shared query with a timeout so one slow query does not hold every waiter indefinitely.
  • Does not keep results after completion (that would be a cache, with invalidation problems).Supporting

The reasoning

  1. Key in-flight queries by every parameter that changes the result.
  2. Remove entries when the query settles, on failure as well as success.
  3. Bound the shared query with a timeout, and don't keep results afterwards.

The code is a dozen lines; the guarantees are in the details. An entry must disappear on failure, the shared work must have a deadline, and the key must capture everything that makes two requests the same. Each of those is an easy omission that turns a protective layer into an outage amplifier.

Stage 9 of 9 · Defend it

Defend the design

Your interviewer pushes back: "This is a lot of machinery. Why not shard Postgres by channel and be done? And isn't the data service just an extra hop that adds latency?"

What you need to know first

"Postgres doesn't scale" isn't true and isn't a defence. A convincing answer names the conditions under which the alternative wins (slower growth, a need for relational queries, more people to run it) and shows which of those conditions doesn't hold here.

"The data service is just an extra hop." What does the hop buy?

Request coalescing and isolation from hot channels, worth far more than a sub-millisecond hop

Without it, a single @everyone ping can saturate three database nodes and slow unrelated channels.

What the stage asks

Answer both questions, saying when you would choose Postgres instead and what the extra hop buys.

Reference answer

Postgres is a fair choice. With slower growth, a team fluent in Postgres, or a need for relational queries across messages, sharding Postgres by channel works well; several large chat products do it. What it costs here is repeated manual resharding: at tens of terabytes a year, you would be splitting shards and moving data again and again, with failover to run yourself. The requirement that capacity grows by adding nodes, for a small team, is what tips it.

The hop is cheap and the protection is not. A request to a data service in the same zone costs well under a millisecond. In exchange, every request for a channel meets in one place, so a burst of identical reads becomes one query, and one hot channel cannot saturate the replicas that also serve thousands of other channels. It is also the one place to put concurrency limits and to dual-write during a migration.

What I am accepting: per-column last-write-wins and the edit/delete race, tombstones that make deletes expensive to read past, compaction and repair to operate, and a query model that only answers "by channel and time". Search, mentions and analytics go to other stores fed from the message stream.

What a strong answer covers

  • Acknowledges sharded Postgres is viable and names when it is the better choice (slower growth, relational queries, existing expertise).
  • Uses the growth rate and team size: resharding by hand repeatedly is the cost being avoided.
  • The data-service hop buys coalescing and isolation from hot keys, worth far more than a sub-millisecond hop.
  • Owns the costs of the chosen design: eventual consistency, tombstones, compaction and repair.

The reasoning

  1. Name when sharded Postgres would win, and why those conditions don't hold here.
  2. The growth rate and team size make repeated manual resharding the cost being avoided.
  3. Own the chosen store's costs: eventual consistency, tombstones, compaction and repair.

The strongest defence names the conditions under which the alternative wins. "Postgres would work if growth were slower or the team bigger" is a much more convincing answer than "Postgres doesn't scale", which is not true.

How you did

Now try it as an interview question

  • “Design Discord.”
  • “Design the message storage for a chat app like Slack or WhatsApp.”
  • “How would you store and paginate billions of chat messages?”
  • “One key in your database is getting far more traffic than the rest. What do you do?”

The interview mode mixes stages from this and other investigations with concept recall and questions about your own projects.

Back to the last stage