Discord: Storing Trillions of Messages (MongoDB → Cassandra → ScyllaDB)
This page is one interview loop in three rounds, built on a real company's public engineering history. All three rounds design the same system: the store that keeps every chat message Discord has ever carried. Each round is an era of that store, and each era ended because the one before it failed in a specific, published way.
| Round 1: Mid-level | Round 2: Senior | Round 3: Architect | |
|---|---|---|---|
| Era | 2015–2016: outgrowing one MongoDB replica set | 2016–early 2022: billions, then trillions, of messages on Cassandra | 2020–2023: data services, ScyllaDB and a migration with no downtime |
| Level (Amazon) | SDE II (L5) | Senior SDE (L6) | Principal (L7) |
| Messages (published) | 100M stored by November 2015; 120M+ sent a day by January 2017 | Billions stored (2017); trillions by early 2022; about 4B sent a day (2022) | Trillions stored, migrated in days |
| Cluster (published) | One MongoDB replica set → 12 Cassandra nodes, 3 copies | 177 Cassandra nodes (early 2022) | 72 ScyllaDB nodes (reported 2023) |
| Traffic we plan for | ~4,200 message writes/s at the peak (from the published daily figure, assumption for the peak) | ~93,000 writes/s at the peak (same method, with a 2× peak factor) | Same writes; 100,000 history fetches/s at the peak (assumption) |
| Target | Stop falling over as data grows | Survive hot channels, mass deletes and a growing cluster | Predictable p99, far less toil, fewer nodes |
| Reading time | ~35 min | ~40 min | ~45 min |
You can start at any round. Rounds 2 and 3 open with a "Where we left off" summary that catches you up.
How to read a case study. Every claim about what Discord actually did comes from Discord's own engineering blog, its public API documentation, or a talk by a Discord engineer, and each round ends with a Sources list. We mark those claims Discord published or (published). Where Discord hasn't published the details, we say so and show a design that fits, labeled as ours. Numbers marked assumption are ours, chosen to make the arithmetic concrete; they are not Discord's internal figures.
Cloud note. Discord runs most of its hardware on Google Cloud, not AWS (Discord, 2022). Where this page shows AWS services, that is a translation for this course, not Discord's setup.
Loop Opener: Why Is Chat History Hard?
A Diary That Never Stops Growing
Discord decided early "to store all chat history forever", so a user can come back years later, on any device, and scroll up. Think of every channel as a diary:
- It only grows. New pages are added at the end, every second, forever. Nothing expires.
- Almost every read is the last page. Opening a channel means "show me the latest 50 messages". Scrolling up means "the 50 before this one".
- Diaries are wildly uneven. Discord described three kinds of servers in 2017: voice-heavy servers that send "a message or two every few days"; private text servers with 100,000 to 1 million messages a year; and large public servers with "thousands of members sending thousands of messages a day".
Synthesizing vector architecture diagram...
Read it as three eras. Each move answered a failure of the store before it: RAM (MongoDB), then hot partitions, garbage-collection pauses and maintenance toil (Cassandra).
A few words we'll use all page:
| Word | What it means on this page |
|---|---|
| Partition key | The part of a row's key that decides which machines store it. All rows with the same partition key live together, on the same replicas. |
| Clustering key | The part of the key that sorts rows inside a partition. |
| Wide partition | A partition holding very many rows, for example years of one channel's messages. |
| Tombstone | A marker that says "this was deleted". In Cassandra-style stores a delete is a write; the data disappears only later. |
| Compaction | The background merge of on-disk files that removes overwritten data and, eventually, tombstones. |
| GC pause | A pause while the Java runtime's garbage collector reclaims memory; during a "stop-the-world" pause, the process answers nothing. |
| Snowflake | A 64-bit ID whose top bits are a timestamp, so IDs sort by creation time. |
| Replica / RF | A copy of a partition on another node; RF (replication factor) is how many copies. |
What Makes It Hard
- Reads are random. Discord found its reads were "extremely random" with a read/write ratio of "about 50/50" (2017). A quiet channel's 50 messages can be spread over months of disk.
- Traffic is concentrated. One announcement in a server with hundreds of thousands of members sends a crowd to the same few rows at the same moment.
- Deletes are expensive. A moderator deleting a million spam messages can hurt a database more than writing them did.
The Question the Whole Loop Answers
How do we store an ever-growing, very uneven pile of messages so that "latest in this channel" is always fast?
The answer grows every round:
- Round 1: match the store to the access pattern: a wide-column database, one partition per channel per 10-day time bucket, rows sorted newest first by Snowflake ID.
- Round 2: learn what that layout costs at scale: hot partitions, tombstones from deletes, compaction backlogs and GC pauses; tune consistency and repairs, and count the toil.
- Round 3: put a routing, request-coalescing data service in front of the database, move to ScyllaDB (no garbage collector, one shard per core), and migrate trillions of messages without downtime.
Round 1 · Mid-level · "Era 1: Outgrowing MongoDB"
~35 min · SDE II (L5) · 2015–2016 · 100M stored messages and growing fast · ~4,200 message writes/s at the peak (planning figure) · one replica set → a 12-node cluster · stop falling over as data grows
R1.1 Establish Design Scope
The interviewer sets the scene: "It's late 2015. We built the first version of our chat app in under two months. Everything, messages included, lives in one MongoDB replica set. We just passed 100 million stored messages and latency has become unpredictable. Design the message store we move to."
| We ask | Interviewer answers | What it changes in the design |
|---|---|---|
| What is the main read? | Open a channel and see the latest messages; then scroll back 50 at a time. Discord (2017): the large public servers "almost always are requesting messages sent in the last hour". | The store must return "the newest N rows of one channel" in one cheap range read (steps 1.1 and 1.4). |
| How are messages written? | One at a time, as people type. Reads and writes are "about 50/50". | We need cheap writes as much as cheap reads; no read-modify-write of big documents. |
| What order do we need? | Order within a channel. There's no order across channels. | We need IDs that sort by time without a central counter (step 1.2). |
| What broke, exactly? | Discord published it: messages sat in a MongoDB collection with one compound index on channel_id and created_at. Around November 2015, at 100 million stored messages, "the data and the index could no longer fit in RAM and latencies started to become unpredictable". | The fix is a store that spreads data across machines by channel, not a bigger box (step 1.1). |
| Do we keep history forever? | Yes. That was decided early. | Partitions grow forever unless we bound them (step 1.3). |
| Do we cache messages? | Discord's requirement (2017): "We also do not want to have to cache messages in Redis or Memcached." Their API alarm fires when the 95th percentile response time goes above 80 ms. | The database itself must have predictable latency; no cache tier to hide it. |
| Who operates it? | Even in January 2017, Discord had "only 4 backend engineers" and no dedicated DevOps engineers. | Low maintenance is a hard requirement, not a nice-to-have. |
Out of scope for this round: trillions of messages, channels with hundreds of thousands of readers, mass deletes, search.
R1.2 Functional Requirements, Derived Step by Step
| Phrase from the problem | Operation |
|---|---|
| "See the latest messages" | getMessages(channel, limit=50) returns the newest messages, newest first |
| "Scroll back" | getMessages(channel, before=message_id, limit=50) |
| "Jump to a message" (a mention or a pin) | getMessages(channel, around=message_id), and after= to scroll forward from there |
| "Send" | createMessage(channel, content) returns the stored message with its ID |
| "Edit" / "Delete" | editMessage(channel, id, content) and deleteMessage(channel, id) |
Discord's 2017 post adds that more random reads were coming: "view your mentions for the last 30 days then jump to that point in history", pinned messages, and full-text search. Every one of them is a jump to an arbitrary point in a channel.
Not yet: trillions of messages, mega-channels, bulk moderation, search (a separate system).
R1.3 Non-Functional Requirements: the Questions
We name each quality first; the numbers come in R1.7. Discord's own list (2017) is the best guide, and we keep its words where they are precise.
- Linear scalability. "We do not want to reconsider the solution later or manually re-shard the data." Adding capacity must mean adding machines.
- Automatic failover. Losing a node must not page anyone.
- Low maintenance. "We should only have to add more nodes as data grows."
- Predictable performance. The API's p95 alarm is 80 ms, and there is no cache to hide slow reads.
- Not a blob store. Appending to a per-channel blob thousands of times a second would mean constant deserializing and rewriting.
- Proven, open source. "We love trying out new technology, but not too new."
R1.4 The API
This API is illustrative: it is modeled on Discord's public API, whose GET /channels/{channel.id}/messages takes before, after or around (only one at a time) and a limit from 1 to 100, defaulting to 50, and returns messages newest first. The internal paths and fields are ours.
Fetch history
httpGET /v1/channels/175928847299117000/messages?before=175928847299117063&limit=50 HTTP/1.1 Authorization: Bearer <session token>
httpHTTP/1.1 200 OK Content-Type: application/json [ { "id": "175928847299117062", "channel_id": "175928847299117000", "author": { "id": "80351110224678912" }, "content": "gg", "edited_timestamp": null }, { "id": "175928847299117061", "channel_id": "175928847299117000", "author": { "id": "41771983423143937" }, "content": "one more?", "edited_timestamp": null } ]
Send a message
httpPOST /v1/channels/175928847299117000/messages HTTP/1.1 Authorization: Bearer <session token> Content-Type: application/json { "content": "one more?", "nonce": "c1-7f3a92" }
httpHTTP/1.1 200 OK Content-Type: application/json { "id": "175928847299117061", "channel_id": "175928847299117000", "content": "one more?", "nonce": "c1-7f3a92" }
- IDs are strings in JSON. Discord's documentation says Snowflakes are "always returned as strings in the HTTP API to prevent integer overflows in some languages" (a JavaScript number is exact only up to 2^53).
- The message ID is the cursor. Because a Snowflake starts with its timestamp,
before=<id>means "older than this message" and needs no separate page token. Discord's docs point out you can even build a cursor from a time:(timestamp_ms - DISCORD_EPOCH) << 22. - The
nonceis a value the client picks so it can match the server's echo to the message it showed optimistically. (Discord's public API has anoncefield; how Discord uses it for retries internally isn't something we rely on.)
Status codes
| Code | Meaning |
|---|---|
200 OK | Messages returned / message stored |
400 Bad Request | A malformed cursor, or limit outside 1–100 |
403 Forbidden | The user may not read this channel's history |
404 Not Found | Unknown channel or message |
429 Too Many Requests | Rate limited; retry after the given delay |
503 Service Unavailable | The store is overloaded; retry with backoff |
Recap
- Five operations, all scoped to one channel.
- Every read is "a slice of one channel, by message ID"; the message ID doubles as the page cursor.
R1.5 Design Evolution: From One Replica Set to a Partitioned Store
Each step is a problem, your turn to think, the answer, and what it costs us.
Step 1.0: The Baseline
One MongoDB replica set (a primary that takes every write, plus secondaries that copy it) holds every collection. Messages have one compound index on channel_id and created_at.
Synthesizing vector architecture diagram...
One primary holds every message and the whole index. It works while the hot data and the index fit in that machine's RAM.
Discord published that this was intentional: build quickly, but "always with a path to a more robust solution", and it had decided not to use MongoDB sharding, which it considered "complicated to use and not known for stability".
Step 1.1: The Index No Longer Fits in RAM
The problem: 100 million messages. The data and the index have outgrown RAM, so reads go to disk, and latency is now unpredictable. Reads are random: a quiet channel's 50 messages are scattered across the disk, and loading them evicts pages that busy channels need. What would you do?
Primitive: Database Sharding and Partition Keys · Loop: Design a Distributed Key-Value Store (partitioning, replication and quorums from the inside)
Step 1.2: Order Messages Without a Coordinator
The problem: MongoDB sorted by created_at. In the new store, the clustering key must be unique and must sort messages in the order they were sent. Many API servers create messages at the same time.
What would you do?
Primitive: Distributed Unique ID Generators · Loop: Design a Distributed Unique ID Generator (clock steps, worker numbers and the exhaustion date, in depth) · Drill: The NTP Sync That Generated Duplicate Order IDs (answered here: a backward clock step and the fix above; why not UUIDv4 in the wrong answers above)
Step 1.3: One Channel's Partition Grows Forever
The problem: with PRIMARY KEY (channel_id, message_id), each channel is one partition forever. While importing existing messages, the logs fill with warnings about partitions over 100 MB. The database's documentation says it can hold 2 GB partitions.
What would you do?
Where a quiet channel's read goes. A voice-heavy server sends, say, one message every two days (assumption: 0.5 a day). Fifty messages span 100 days, so "the latest 50" walks back through about 11 buckets, one query each, most of them nearly empty. Discord accepted this: "for active Discords enough messages are usually found in the first partition and they are the majority." Round 2 adds a fix.
Synthesizing vector architecture diagram...
A quiet channel's read walks backwards through buckets, one partition query each, until it has 50 messages or reaches the channel's creation bucket. A busy channel stops at the first box.
Step 1.4: Reads Must Be "Latest First"
The problem: a partition is sorted by message_id. Almost every read wants the newest rows first.
What would you do?
Round 1 Step Summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 1.0 | (baseline) | One MongoDB replica set, compound index | Everything must fit one machine's RAM |
| 1.1 | Index outgrew RAM | Cassandra, partitioned by channel (KKV) | New operations skills; one table per query |
| 1.2 | Order without a coordinator | Snowflake IDs, Discord epoch 2015 | Clock handling |
| 1.3 | Partitions grow forever | (channel_id, bucket), 10-day buckets | Reads may span buckets |
| 1.4 | Newest first | CLUSTERING ORDER BY (message_id DESC) | Ascending scans are reverse queries |
R1.6 Architecture v1
Synthesizing vector architecture diagram...
Every request names a channel, so it lands on the three replicas that own that channel's current bucket. Any node can coordinate a request; data spreads over the ring by partition key.
The table (Discord's published minimal schema; the real one has 16 columns):
sqlCREATE TABLE messages ( channel_id bigint, bucket int, message_id bigint, author_id bigint, content text, PRIMARY KEY ((channel_id, bucket), message_id) ) WITH CLUSTERING ORDER BY (message_id DESC);
Synthesizing vector architecture diagram...
One channel owns many small partitions, one per 10-day window; inside each, rows sort by Snowflake, newest first.
How we switched (published). Discord imported existing messages, then "set up our code to double read/write to MongoDB and Cassandra" (a dark launch: the new store gets real traffic while the old one still serves users). After it held up, Cassandra became the primary store and MongoDB was "phased out within a week". Measured in testing: writes "sub-millisecond" and reads "under 5 milliseconds", consistently for a week, and a jump back a full year in a channel with millions of messages stayed fast. The dark launch also surfaced the first surprise (a message with a null author_id), which is Round 2's step 2.4.
Trace: send, then open the channel (our design, with Discord's key layout)
Synthesizing vector architecture diagram...
The write is acknowledged by two of three replicas; the third catches up in the background. The read touches one partition when the channel is busy.
The consistency level (LOCAL_QUORUM, two of three copies) is our choice for this round; Discord published in 2023 that it performs "reads and writes with quorum consistency level", which Round 2 relies on.
Sources for this round
- Vishnevskiy, How Discord Stores Billions of Messages, Discord blog, January 2017: store history forever; MongoDB replica set and compound index; 100M messages and RAM in November 2015; the read patterns and server types; the requirements; the choice of Cassandra; KKV;
(channel_id, message_id); 100 MB partitions; 10-day buckets and the bucket function; the dark launch; measured latencies; 12 nodes, RF 3, nearly 1 TB per node; 40M, 100M and 120M+ messages a day. - Discord Developer Documentation, API Reference: Snowflakes and Message resource: Get Channel Messages (checked September 2026): the ID layout, the epoch, IDs as strings,
before/after/aroundandlimit. - Ingram, How Discord Stores Trillions of Messages, Discord blog, March 2023: quorum reads and writes; the minimal schema; reverse queries.
Everything else in this round (the peak factor, row sizes, the API shapes beyond the public fields, the consistency level in the trace) is our design or our assumption.
R1.7 Numbers
Figures marked published come from the sources above; the rest are assumptions or derived from them.
Write and read rates
| Quantity | Arithmetic | Result |
|---|---|---|
| Messages a day | Published, January 2017 | 120,000,000 |
| Average writes | 120,000,000 ÷ 86,400 s | 1,389/s |
| Peak writes | 3 × average (assumption: evening peaks) = 4,167 | ≈ 4,200/s |
| Peak history reads | Read/write "about 50/50" (published) | ≈ 4,200/s |
Cluster load at the peak (12 nodes, RF 3; reads at LOCAL_QUORUM, our choice)
| Quantity | Arithmetic | Per node |
|---|---|---|
| Replica writes | 4,200 × 3 copies = 12,600/s ÷ 12 nodes | 1,050/s |
| Replica reads | 4,200 × 2 replicas × 1.3 buckets per read (assumption: some reads walk back) = 10,920/s ÷ 12 | 910/s |
Both are small for a Cassandra node. Round 1's problem is data size and layout, not request rate.
Storage: a sanity check against the published cluster. Discord reported "nearly 1TB of compressed data on each node" on 12 nodes: about 12 TB, or 12 ÷ 3 = 4 TB of unique data. How many messages is that? Discord didn't say. Draw daily volume piecewise through the published points (about 10M a day in January 2016, our assumption; 40M a day in July 2016; 100M in December 2016; 120M in January 2017) and add up the days: 2016 held roughly 20 billion messages. Then 4 TB ÷ 20B ≈ 200 bytes per message, compressed, all overhead included. That's our estimate, not a published figure; it only says the published numbers hang together.
Growth. At 120M messages a day × 200 B = 24 GB of new unique data a day; × 3 copies = 72 GB a day across the cluster, or 72 ÷ 12 = 6.0 GB per node a day. A node gains another terabyte in about 1,000 ÷ 6.0 ≈ 167 days, sooner as traffic grows. Discord's plan: newer Cassandra versions handle more data per node, so it expected to go from 1 TB to 2 TB per node, and Cassandra 3's storage format "can reduce storage size by more than 50%". Then keep adding nodes.
Buckets. From step 1.3: 10-day buckets hold up to 200,000 rows at 500 B before reaching 100 MB. The busy channel at 5,000 messages a day uses 25 MB per bucket; the quiet channel at 0.5 a day holds 5 messages per bucket and walks about 11 buckets for 50 messages.
Synthesizing vector architecture diagram...
A 10-day bucket keeps even a very busy channel at or under 100 MB, while a quiet channel's buckets are almost empty: that's the price we pay on reads.
R1.8 Trade-Offs
MongoDB (one replica set) vs Cassandra for this access pattern
| MongoDB replica set | Cassandra (chosen) | |
|---|---|---|
| Writes | One primary takes all of them | Any node coordinates; each write goes to 3 replicas |
| Growth | A bigger machine, or sharding Discord judged too risky | Add nodes; data spreads by partition key |
| "Latest 50 in a channel" | Index walk plus scattered document reads | One partition, adjacent rows on disk |
| Failure | Primary election pauses writes | Losing one of three replicas is invisible at quorum |
| Queries | Flexible | Only what the keys allow |
Bucket size
| Smaller buckets (1 day) | 10 days (Discord) | Larger buckets (100 days) | |
|---|---|---|---|
| Busiest partition | 2.5 MB | 25 MB | 250 MB: past the warning line |
| Quiet channel, latest 50 | ~101 queries (0.5/day) | ~11 queries | ~2 queries |
| Spread of a busy channel | Moves to new replicas daily | Every 10 days | Stuck on 3 replicas for months |
The right bucket is set by the busiest channel (size) and paid for by the quietest (queries). One global size is a compromise; Round 2 reduces the quiet channels' cost instead of changing the size.
Row per message vs a blob per channel. A blob per channel (one document holding many messages) would mean deserializing and rewriting the blob on every send, thousands of times a second: Discord's "not a blob store" requirement.
R1.9 Failure Modes
| Trigger | What you'd see | How the design responds |
|---|---|---|
| A node dies | One of three replicas missing for about a quarter of partitions (with 12 nodes, each node holds 3/12 of the data) | Quorum needs 2 of 3, so reads and writes continue. The coordinator stores hints (writes for the dead node) and replays them when it returns, within a window (Cassandra's default max_hint_window_in_ms is 3 hours); after longer, a repair (an anti-entropy comparison between replicas) fixes it. |
| A node is slow | Some requests wait on it | A quorum read needs only 2 replicas; drivers and coordinators can retry a slow replica read elsewhere (speculative retry). A slow coordinator still hurts every request it coordinates: Round 2's GC story. |
| A clock steps backwards on an API server | Its IDs could repeat or sort before older messages | The generator refuses to issue a timestamp below the last one it used and waits it out; alarm on clock offset. |
| The dark launch finds a bug | Errors in the bug tracker, not user impact | The old store still serves reads; fix and continue. Discord's author_id null surprise was found exactly this way. |
| A bucket boundary passes | Nothing visible | New messages go to the next bucket automatically; reads span both until the new one has 50 messages. |
R1.10 Pillar Check
| Pillar | What Round 1 covers |
|---|---|
| Reliability | Three copies of every partition and quorum reads and writes; a node loss is invisible; the move to Cassandra ran as a dark launch with both stores written REL 11 · REL 8 |
| Performance Efficiency | The key design follows the one query that matters: one partition, adjacent rows, newest first; buckets bound partition size PERF 3 |
| Security | Every history read is checked against the user's right to read the channel before the store is touched; IDs travel as strings so clients don't corrupt them SEC 3 |
| Cost Optimization | An open-source store that scales by adding commodity nodes, and no cache tier to pay for COST 5 |
| Operational Excellence | Skipped this round: repairs and node operations become Round 2's main story. |
| Sustainability | Skipped this round: at 12 nodes, the footprint isn't yet the question. |
R1.11 Round 1 Rubric and Follow-Ups
What a strong mid-level (L5) answer shows
- Starts from the access pattern ("latest 50 in one channel") and picks a store whose layout matches it.
- Designs the key: partition by channel, cluster by a time-sortable ID, newest first.
- Knows why
created_atcan't be the clustering key and why random UUIDs can't be the cursor. - Bounds partition growth with time buckets, and can derive the bucket from the ID.
- Names what buckets cost (multi-partition reads for quiet channels) and accepts it on purpose.
- Keeps the old store running while the new one proves itself.
Follow-up questions
-
"Why not keep MongoDB and shard it?" Answer: that was an option Discord explicitly rejected at the time as "complicated to use and not known for stability". The deeper reason: even sharded, the data layout (documents plus a secondary index) doesn't put a channel's recent messages next to each other on disk the way a clustered partition does.
-
"The bucket size is global. A single channel suddenly gets 50,000 messages a day. What happens?" Answer: that channel's buckets grow to 50,000 × 10 × 500 B = 250 MB, over the target. Nothing breaks immediately, but the partition becomes slower to compact and read. Options: a per-channel bucket size (stored with the channel, so reads can compute it), or a smaller global size. Discord kept one size and fixed the quiet-channel cost instead.
-
"When does a Discord Snowflake run out?" Answer: 42 bits of milliseconds is 2^42 ms ≈ 139 years after 2015, around 2154. But stored in a signed 64-bit
bigint, the top bit becomes the sign bit after 2^41 ms ≈ 69.7 years, around 2084, when new IDs would turn negative and sort below old ones. (That's our observation from the published layout, not something Discord has written about.)
Interview gotchas from this round's wrong answers
| Gotcha | Why it's wrong |
|---|---|
| "A bigger machine" | History grows forever; it only delays the same wall. |
| "Cache the messages" | Quiet channels' random reads are exactly what a cache misses, and Discord ruled it out. |
"created_at as the clustering key" | Two messages in the same millisecond collide and one overwrites the other. |
| "Partition by channel; the database allows 2 GB" | Big partitions pressure compaction and memory and pin a channel to three replicas forever. |
| "Sort in the application" | Reads a whole bucket to return 50 rows. |
Round 2 · Senior · "Era 2: Billions of Messages on Cassandra"
~40 min · Senior SDE (L6) · 2016–early 2022 · 12 → 177 Cassandra nodes (published) · billions, then trillions, of messages · about 4B messages sent a day (published, 2022) · ~93,000 writes/s at the peak (planning figure) · survive hot channels, mass deletes and the cluster's own maintenance
R2.0 Where We Left Off
This is what the candidate says aloud in the first 60 seconds of Round 2. If you're starting here, it's everything you need from Round 1.
Round 1 in 60 seconds. "Discord began with every message in one MongoDB replica set, indexed on channel and creation time. At 100 million messages in late 2015, the data and index outgrew RAM and latency became unpredictable. Reads are random and about half of all traffic, and Discord didn't want a cache in front, so the store itself had to match the access pattern. It moved to Cassandra: the partition key is
(channel_id, bucket), where a bucket is a 10-day window derived from the message's Snowflake ID, and rows are clustered bymessage_idnewest first. 'Latest 50' reads one partition for a busy channel and walks back through buckets for a quiet one. Snowflakes are 64-bit IDs with a millisecond timestamp since 2015, so they sort by time and work as page cursors. The move ran as a dark launch, writing both stores. By January 2017 it was 12 nodes with 3 copies of everything. Open costs: a very busy channel still lives on three replicas at a time, deletes are writes, and the cluster's maintenance is ours to run."
Architecture v1, compact
Synthesizing vector architecture diagram...
Round 1 in one picture: every read is one channel's slice, newest first, from three replicas.
Round 1 step summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 1.1 | Index outgrew RAM | Cassandra, partitioned by channel | One table per query |
| 1.2 | Order without a coordinator | Snowflake IDs | Clock handling |
| 1.3 | Partitions grow forever | 10-day buckets | Reads may span buckets |
| 1.4 | Newest first | Descending clustering order | Ascending scans are reverse queries |
Open costs: hot channels land on three replicas; deletes are writes; repairs, compaction and node operations are ours; eventual consistency is the default.
R2.1 The Scope Raise
Interviewer: "It's a few years later. We have hundreds of thousands of people in some servers, and when one of them posts an announcement, a crowd opens the same channel at once. Moderators delete spam by the thousands. The cluster has grown from 12 nodes to well over a hundred, and our on-call engineers are firefighting it every week. Keep reads fast and make it survivable."
| We ask | Interviewer answers | What it changes in the design |
|---|---|---|
| How big is the cluster now? | 177 nodes at the beginning of 2022, holding trillions of messages; the average node has about 4 TB of disk. (Discord, 2023.) | About 177 × 4 = 708 TB of disk; 15 times as many nodes to maintain as in 2017. |
| How many messages a day? | About 4 billion a day (Discord, 2022). | ~46,300 writes/s on average; we plan for ~93,000/s at the peak (R2.6). |
| What does a busy channel look like? | "A server with just a small group of friends tends to send orders of magnitude fewer messages than a server with hundreds of thousands of people." An @everyone announcement sends a crowd to one channel. | One partition takes a crowd's reads: a hot partition (step 2.1). |
| What happens when moderators delete? | Deletes arrive in bulk. Discord's public API today deletes 2 to 100 messages per call, only messages under 2 weeks old; in 2016 a server deleted millions of messages through the API. | Deletes create tombstones that reads must scan past (step 2.2). |
| What consistency do we use? | Reads and writes at quorum. (Discord, 2023.) | A slow replica slows every quorum request that needs it (steps 2.1, 2.4). |
| What does on-call look like? | "Our on-call team was frequently paged for issues with the database, latency was unpredictable, and we were having to cut down on maintenance operations that became too expensive to run." | Toil becomes a design input, and the case for the next era (step 2.5). |
Scope change
| Round 1 | Round 2 | |
|---|---|---|
| Messages stored | Billions | Trillions |
| Messages a day | 120M | ~4B |
| Nodes | 12 | 177 |
| Busiest channel | Thousands of members | Hundreds of thousands of members, crowds on one announcement |
| Deletes | Occasional | Bulk moderation; mass deletes |
| Operations | "It should just work" | Frequent pages; maintenance being cut back |
The "Not yet" list from R1.2 comes back: mega-channels and mass deletes are in scope now.
R2.2 What Breaks in the Round 1 Design
| Round 1 piece | What breaks at the new scale |
|---|---|
| One partition per channel per 10 days | A crowd reads one partition. Discord: "Lots of concurrent reads as users interact with servers can hotspot a partition." The three replicas that own it fall behind, and at quorum every other request to those nodes waits too. |
| Deletes as writes | Every deleted message leaves a tombstone that reads of that partition must skip until compaction removes it. A mass delete turns "latest 50" into a scan of tombstones. |
| A Java database | Garbage-collection pauses stop a node; Discord "spent a large amount of time tuning the JVM's garbage collector and heap settings, because GC pauses would cause significant latency spikes". |
| Background compaction | "We were prone to falling behind on compactions": reads touch more files, and a node trying to catch up adds "cascading latency". |
| Repairs and node operations | They grow with the data. Discord had to "cut down on maintenance operations that became too expensive to run". |
R2.3 New Requirements and API Additions
Bulk delete (the shape of Discord's public endpoint; the internal behavior is ours):
httpPOST /v1/channels/175928847299117000/messages/bulk-delete HTTP/1.1 Authorization: Bearer <moderator session token> Content-Type: application/json { "messages": ["175928847299117061", "175928847299117062"] }
httpHTTP/1.1 204 No Content
Discord's public endpoint requires the MANAGE_MESSAGES permission, takes 2 to 100 IDs, and refuses messages older than 2 weeks with 400 Bad Request. That last rule is a product limit worth noticing: it keeps mass deletes inside the newest buckets.
Edit
httpPATCH /v1/channels/175928847299117000/messages/175928847299117061 HTTP/1.1 Content-Type: application/json { "content": "one more game?" }
The store writes only the changed columns (content, edited_timestamp). In Cassandra every write is an upsert (insert-or-update, with no read first), which matters in step 2.4.
Edit history (our design; Discord hasn't published whether or how it keeps one): each edit overwrites content, so keeping earlier versions would mean a separate message_edits table keyed by (channel_id, bucket) like messages and clustered by (message_id, edit_timestamp), written alongside the edit and deleted with the message.
Consistency per call (our design, around Discord's published quorum choice)
| Operation | Write or read level | Why |
|---|---|---|
| Send, edit, delete | QUORUM (2 of 3) | Survives one replica down; a quorum read afterwards sees it |
| Fetch history | QUORUM | Overlaps every quorum write, so an acknowledged message is never missing from the next fetch |
| Background jobs (exports, analytics) | ONE | Speed over freshness; a repair or the next run fixes gaps |
Read repair, in one line: when the replicas a quorum read contacts disagree, the coordinator writes the newest version to the stale replica before answering (Cassandra calls this blocking read repair), so a read also heals what it touched.
R2.4 Design Evolution: Hot Partitions, Tombstones, Pauses and Toil
Step 2.1: A Viral Channel Makes One Partition Hot
The problem: a server with 500,000 members posts @everyone. Within a minute, 100,000 people open the channel (assumption). All of them read the newest bucket of one channel: one partition, on three nodes. Latency on those nodes climbs, and soon unrelated channels on the same nodes are slow too. What would you do?
Primitive: Database Sharding and Partition Keys · Loop: Design a Chat and Instant Messaging System (gateways, fan-out and channels with 150,000 members; this page covers only the store)
Step 2.2: Deleting Spam Slowed Reads to 20 Seconds
The problem (published; about mid-2016, derived from the post's dates): about six months after the switch, Cassandra ran "10 second 'stop-the-world' GC constantly". The team traced it to one public server's channel that took 20 seconds to load. The channel had one message in it. The server had deleted millions of messages through the API. What would you do?
Go deeper: the 2-day timing chain. Everything that can delay a delete reaching a replica has to fit inside the grace period.
| Time | Event |
|---|---|
| T0 − 1 h | Replica C goes down. Coordinators start keeping hints for it. |
| T0 | A moderator deletes a message. A and B write the tombstone; the coordinator stores a hint for C. |
| T0 + 2 h | Three hours after C went down, coordinators stop writing new hints for it (Cassandra's default max_hint_window_in_ms). Writes from now on reach C only through repair. |
| Night 1 and night 2 | Nightly repairs run, but C is down, so its ranges can't be repaired. |
| T0 + 48 h | The tombstone on A and B passes gc_grace_seconds (2 days). Any compaction may now drop it. |
| T0 + 50 h | C comes back with the message still live. The hint from T0 can't save us: Cassandra gives a hint a lifetime no longer than the table's gc_grace_seconds, for exactly this reason, so it expired at T0 + 48 h. The next repair sees "live" on C and "nothing" on A and B, and copies the deleted message back. |
The rule that follows: a node that has been down longer than the grace period must not rejoin with its old data. Wipe it and rebuild it from the others; that's the standard operating rule for Cassandra-style stores. A replica is a safe repair source only if it has seen every delete up to the moment the others purged theirs. At 10 days, that rule rarely bites; at 2 days, a long weekend outage triggers it.
Primitive: Write-Ahead Log and LSM Trees · Loop primitive: LSM-Trees & Compaction (Part 2 on versions and tombstones, and Part 5's worked example of a tombstone that must stay)
Step 2.3: Latency Spikes of Seconds
The problem: on a normal afternoon, p99 read latency on some nodes jumps from tens of milliseconds to seconds, for a few seconds at a time. Nothing in the traffic explains it. The same nodes are behind on compaction. What would you do?
Primitive: Write-Ahead Log and LSM Trees · Loop primitive: LSM-Trees & Compaction (Part 5, "Why writes stall", with the numbers) · Drill: The Time-Series Database That Froze on Flush (both questions answered in this step)
Step 2.4: A Message With No Author
The problem (published): right after the dark launch, the bug tracker filled with errors: author_id is null. It's a required field. Investigation: a user edited a message at the same moment another user deleted it.
What would you do?
Primitive: Distributed Consensus: Raft and Paxos · Loop: Design a Distributed Key-Value Store (Round 2: tunable consistency, collisions, conditional writes, hinted handoff and anti-entropy)
Step 2.5: Operations Eat the Team
The problem: by early 2022 the cluster has 177 nodes. On-call is paged often, weekends go to firefighting, and repairs and other maintenance are being cut back because they cost too much to run. What would you do?
Round 2 Step Summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 2.1 | A viral channel's partition is hot | Understand the key; rate limits, push-first history, speculative retry | Unbounded concurrency remains → Round 3 |
| 2.2 | Deletes slowed reads | Tombstones understood; gc_grace 10 → 2 days with nightly repair; empty-bucket tracking; no null writes | Tight repair schedule |
| 2.3 | Latency spikes of seconds | GC pauses and compaction backlog; gossip dance | Constant tuning toil |
| 2.4 | A message with no author | LWW per column understood; detect-and-delete on read (or LWT) | A check on read, or Paxos per edit |
| 2.5 | Operations eat the team | Measure the toil | The case for Round 3 |
R2.5 Architecture v2
Synthesizing vector architecture diagram...
The data path barely changed from Round 1; the operations box grew. The gateway (see the chat loop) delivers new messages; the store serves history.
Trace: a hot channel
Synthesizing vector architecture diagram...
Nothing limits how many identical reads reach the three replicas at once. Other channels whose data shares those nodes suffer too.
Trace: a bulk delete, then a read
Synthesizing vector architecture diagram...
A mass delete makes the next reads more expensive, not cheaper, until compaction catches up after the grace period.
Sources for this round
- Vishnevskiy, How Discord Stores Billions of Messages, Discord blog, January 2017: the dark launch and the
author_idrace; tombstones, 10-day default, 12 null tombstones per message; the 2016 incident (10 s GC, 20 s load, millions deleted); 2-day grace with nightly repairs; empty-bucket tracking; interest in Scylla for repair times. - Ingram, How Discord Stores Trillions of Messages, Discord blog, March 2023: 177 nodes early 2022; hot partitions and quorum; compaction backlog and the gossip dance; GC tuning and reboots; p99 figures on Cassandra; other clusters' faults.
- Oakley, How Discord Supercharges Network Disks for Extreme Low Latency, Discord blog, August 2022: 4 billion messages a day; Google Cloud.
- Discord Developer Documentation, Bulk Delete Messages (checked September 2026).
- Apache Cassandra documentation, cassandra.apache.org/doc: tombstones and
gc_grace_seconds, hinted handoff, compaction strategies, lightweight transactions and serial consistency.
Discord didn't publish its interim hot-partition measures, its consistency choices per job, or the toil figures; those are ours.
R2.6 Numbers and Cost
All figures are assumptions unless marked published.
Writes
| Quantity | Arithmetic | Result |
|---|---|---|
| Messages a day | Published, 2022 | ~4,000,000,000 |
| Average writes | 4,000,000,000 ÷ 86,400 s | 46,296/s |
| Peak writes | 2 × average (assumption; a global user base flattens the peak more than Round 1's) | ≈ 93,000/s |
| Replica writes per node at peak | 93,000 × 3 = 279,000/s ÷ 177 nodes | ≈ 1,576/s |
Raw write rate isn't what hurt. The pain was concentration (hot partitions), deletes and the runtime.
The cluster
| Quantity | Arithmetic | Result |
|---|---|---|
| Nodes | Published, early 2022 | 177 |
| Disk | 177 × ~4 TB average per node (published) | ≈ 708 TB of disk |
| Growth since 2017 | 177 ÷ 12 | ~15× the nodes |
One hot channel (from step 2.1)
| Quantity | Arithmetic | Result |
|---|---|---|
| Opens | 100,000 in 60 s (assumption) | 1,667 reads/s on one partition |
| Replica reads | × 2 (quorum) | 3,333/s on 3 nodes, ≈ 1,111 per node for one key |
| In flight (Little's law) at 5 ms / 100 ms / 1 s | 1,667 × 0.005 / 0.1 / 1 | 8 / 167 / 1,667 |
Tombstones from a bulk delete (illustrative)
| Quantity | Arithmetic | Result |
|---|---|---|
| Messages deleted | 10,000 calls × 100 IDs | 1,000,000 |
| Tombstone bytes | 1,000,000 × ~30 B (assumption: key plus deletion time) | ≈ 30 MB per replica, 90 MB for 3 replicas |
| How long they cost reads | gc_grace_seconds + time to the next compaction that includes them | ≥ 2 days at Discord's setting (172,800 s), ≥ 10 days at the default (864,000 s) |
| Null tombstones avoided | 120M messages a day (2017) × 12 | 1.44 billion a day |
GC pauses and p99 (illustrative)
| Paused share of time per node | Pause length | Effect |
|---|---|---|
| 0.1% | 1 s | 0.3% meet a pause: p99 unaffected, p99.9 ≈ 0.67 s |
| 1% | 1 s | 3% meet a pause: p99 ≈ 0.67 s, p99.9 ≈ 0.97 s |
| 1% | 10 s (Discord's 2016 incident) | p99 ≈ 6.7 s, and clients time out |
The cost that matters here is effort (COST 11). The machines were affordable; the people weren't. Step 2.5's table is the bill: pages, dances, reboots and repairs that grow with every node.
R2.7 Trade-Offs
Consistency levels (RF 3)
| Write / read | Survives | Read sees the latest write? | Latency |
|---|---|---|---|
ONE / ONE | 2 replicas down | No: may read a replica that missed it | Lowest |
QUORUM / QUORUM (Discord) | 1 replica down | Yes: 2 + 2 > 3, the sets overlap | The slower of 2 replicas |
ALL / ONE | No replica down for writes | Yes | Writes wait for the slowest replica |
LWT (SERIAL or LOCAL_SERIAL) | 1 replica down | Yes, and conditional | Several round trips |
Quorum's price is visible in step 2.1: when a replica is slow, a quorum request that needs it is slow. ONE would hide a slow replica but lets a fetch miss an acknowledged message.
Bucket size, revisited. Changing the bucket size now means rewriting trillions of rows into new partitions, which is a migration. Tracking empty buckets fixed the quiet-channel cost without touching the layout.
Soft delete vs hard delete
| Hard delete (tombstone) | Soft delete (a deleted flag, content cleared) | |
|---|---|---|
| Read cost | Tombstones scanned until purged | A row per deleted message, forever, skipped in the app |
| Storage | Reclaimed after grace and compaction | Row stays (small once content is cleared) |
| Resurrection risk | Yes, if a replica misses it past the grace period | No tombstone to lose; the flag is an ordinary write |
| Data really removed | Yes, eventually | Only the cleared columns; clearing them is itself a tombstone per column |
A soft delete doesn't escape tombstones: clearing content writes one. It only trades a row tombstone for a cell tombstone plus a permanent row. Discord kept hard deletes and fixed the read path.
R2.8 Failure Modes
| Trigger | What you'd see | How the design responds |
|---|---|---|
| Hot partition overload | p99 up on 3 nodes, then on everything that shares them | Rate limits per channel; history fetch starts from the client's newest pushed message; speculative retry; page on per-node p99, not the cluster average. The real fix is Round 3's concurrency control. |
| Tombstone accumulation | Warnings in logs for queries over 1,000 tombstones; slow loads for one channel | Skip empty buckets; compaction after the grace period; alarm on tombstones scanned per read. |
| GC pause storms | Seconds-long stalls; a node flaps up and down | Tune heap and collector; take the node out of rotation; reboot if pauses chain. This is toil, not a fix. |
| Compaction backlog | Pending compactions grow; SSTables per read grow; latency cascades | The gossip dance, inside the 3-hour hint window; add nodes; lower write amplification. |
A node down longer than gc_grace_seconds | It holds deleted messages as live | Never rejoin it with its data: wipe and rebuild from the other replicas (step 2.2). |
| An edit races a delete | A row with no author_id | Treat as deleted on read and remove it (step 2.4). |
| A retried write lands after a newer one | An edit timed out and was retried; the original, delayed, reaches a coordinator later than a second, newer edit | Last write wins by timestamp. If the timestamp is assigned by the coordinator on arrival, the delayed original gets a later timestamp and wins. Fix: set the write timestamp once per logical edit on the client side (USING TIMESTAMP, or the driver's client-side timestamps) and reuse it on every retry. |
R2.9 Production Gotchas
| Gotcha | Why it hurts | What we do |
|---|---|---|
| Unbounded wide partitions | GC pressure in compaction and streaming; a partition can't spread | Bucket by time from day one; alarm on partition size (Cassandra versions of that era warned at 100 MB by default, which is what Discord saw) |
| Hard deletes at volume | Tombstones cost reads until purged | Product limits (Discord's bulk delete: 100 per call, under 2 weeks old); skip empty buckets; alarm on tombstones per read |
| Writing nulls | Every null is a tombstone | Write only the columns that have values |
| Relying on the OS page cache for hot data | In 2017, large public servers' recent messages were "usually in the disk cache", which the OS can evict at any time; quiet channels' reads cause "disk cache evictions" for everyone | Treat the page cache as luck, not capacity; Round 3's store keeps its own cache and is shielded by coalescing |
Lowering gc_grace_seconds without a repair schedule | Deleted messages come back | Repairs must complete within the grace period, every time |
| Timestamps assigned on arrival | A delayed retry beats a newer write | Client-side timestamps, reused across retries |
R2.10 Pillar Check
| Pillar | What Round 2 adds |
|---|---|
| Reliability | Quorum reads and writes; hints and nightly repair inside the grace period; nodes down past the grace period are rebuilt, not rejoined REL 11 |
| Performance Efficiency | Understands why one key concentrates load; skips empty buckets; writes no nulls; knows p99 is set by the worst node PERF 3 · PERF 5 |
| Security | Deletes that really remove data (tombstone, grace, compaction), and don't come back; moderation actions require a specific permission SEC 3 |
| Cost Optimization | The cost that grew was engineering effort, measured in pages, dances and repairs COST 11 |
| Operational Excellence | Alarms per node on p99, pending compactions, tombstones per read and partition size; a written procedure for the gossip dance OPS 8 · OPS 10 |
| Sustainability | Skipped this round: the node count is about to fall by more than half, in Round 3. |
R2.11 Round 2 Rubric and Follow-Ups
What a senior (L6) answer adds over L5
- Explains a hot partition from first principles: one key, three replicas, reads costlier than writes, and quorum spreading the pain.
- Uses Little's law to show how a slow node turns into a queue.
- Knows deletes are tombstones, why the grace period exists (so repair can't resurrect deletes), and the rule for a node down past it.
- Connects GC pauses to p99 with arithmetic, and compaction backlog to GC.
- Knows last-write-wins is per column, and the LWT alternative with its cost and its
SERIAL/LOCAL_SERIALchoice. - Treats operator toil as a cost to measure.
Follow-up questions
-
"Why did Discord lower the grace period instead of raising it?" Answer: a shorter grace period purges tombstones sooner, so mass deletes stop costing reads sooner. It's safe only because repairs ran every night, well inside 2 days. The price is the rule that a node down for more than 2 days must be rebuilt.
-
"Would
LOCAL_ONEreads fix hot partitions?" Answer: they'd halve the replica reads (one instead of two) and hide one slow replica, but a crowd still hits the same three nodes, and a fetch may miss a just-acknowledged message. It's a 2× trade, not a fix. -
"Why not TTL old messages away?" Answer: the product promise is history forever. Expiring rows would also create tombstones (an expired cell behaves like one until compaction).
Interview gotchas from this round's wrong answers
| Gotcha | Why it's wrong |
|---|---|
| "Add nodes to fix a hot partition" | One key stays on the same replicas. |
| "Deletes free space immediately" | They add tombstones that cost reads for at least the grace period. |
"Shorten gc_grace to purge faster" | Only safe if repair always completes inside it. |
| "Tune GC until it stops" | Discord tuned for years; pauses still dominate the p99. |
| "Quorum orders my edits and deletes" | It doesn't; last write wins per column. |
| "The OS page cache is our cache" | It's shared, evictable and invisible to capacity planning. |
Round 3 · Architect · "Era 3: Trillions, ScyllaDB and Data Services"
~45 min · Principal (L7) · 2020–2023 · trillions of messages on 177 Cassandra nodes → 72 ScyllaDB nodes (published) · ~93,000 writes/s and 100,000 history fetches/s at the peak (planning figures) · migrate everything with no downtime · predictable p99, far less toil
R3.0 Where We Left Off
Round 2 in 60 seconds. "The layout from Round 1 held:
(channel_id, bucket)partitions of 10 days, newest first. But at 177 nodes and trillions of messages, Cassandra became a high-toil system. A crowd opening one channel makes its partition hot; reads cost more than writes, and with quorum reads and writes, every request to those three nodes slows down. Deletes are tombstones that reads scan past until compaction purges them after the grace period; Discord cut the grace period from 10 days to 2 because repairs ran nightly, skipped empty buckets, and stopped writing nulls, which had been 12 tombstones per message. GC pauses dominated the p99, compaction fell behind, and the 'gossip dance' took nodes out to compact. On Cassandra, history fetches ran at a p99 of 40 to 125 ms and inserts 5 to 70 ms. Open costs: nothing limits concurrency on a hot partition, the runtime pauses, and maintenance grows with every node."
Architecture v2, compact
Synthesizing vector architecture diagram...
Round 2 in one picture: the same data path as Round 1, with the operations box now the biggest thing on the page.
Round 2 step summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 2.1 | Hot partitions | Key analysis; rate limits; push-first history | Unbounded concurrency remains |
| 2.2 | Deletes slowed reads | Tombstones; 2-day grace with nightly repair; skip empty buckets | Tight repair schedule |
| 2.3 | Latency spikes | GC and compaction understood; gossip dance | Toil |
| 2.4 | Edit races delete | LWW per column; detect on read | A check on read |
| 2.5 | Toil | Measured | The case for this round |
Open costs: unbounded concurrency on hot partitions, a garbage-collected runtime, maintenance that no longer fits its windows.
R3.1 The Scope Raise
Interviewer: "It's early 2022. Every one of our other databases has moved to ScyllaDB. The messages cluster is the last one on Cassandra: nearly 200 nodes and trillions of messages. When a big server posts an announcement to everyone, thousands of clients ask for the same messages at the same moment. We want predictable p99s and fewer nodes, and we want to move trillions of messages without any downtime, quickly, because we're firefighting the old cluster."
| We ask | Interviewer answers | What it changes in the design |
|---|---|---|
| Why ScyllaDB? | Discord (2023): "better performance, faster repairs, stronger workload isolation via its shard-per-core architecture, and a garbage collection-free life". By 2020 it had "migrated every database but one to ScyllaDB". | We have production experience; the question is how to move the biggest cluster (step 3.1). |
| Will a new database fix hot partitions? | Discord's own answer: "Hot partitions can still be a thing in ScyllaDB". | We protect the database from upstream: a data service tier (step 3.2). |
| What does a hot moment look like? | "A big announcement on a large server that notifies @everyone: users are going to open the app and read the message." | Identical concurrent reads must be merged (step 3.2). |
| How fast must the migration be? | Discord: "with no downtime, and we need to do it quickly". | Dual writes plus a fast backfill, with verification (step 3.3). |
| Where does it run? | Google Cloud. GCP's local NVMe SSDs are fast but Discord found them not reliable enough on their own for critical data; network persistent disks are durable, but Discord describes them as disks "that take a millisecond or two to complete an operation". | A storage layout that reads locally and writes durably (step 3.1, the super-disk). |
| Can we archive old messages? | Discord floated archiving unused channels to Google Cloud Storage in 2017, and wrote: "We want to avoid doing this one and don't think we will have to do it." | We evaluate it, and probably don't build it (step 3.5). |
Scope change
| Round 2 | Round 3 | |
|---|---|---|
| Database | Cassandra (Java) | ScyllaDB (C++), Cassandra-compatible |
| Between API and database | Nothing | Rust data services: routing and request coalescing |
| Hot partitions | Unbounded concurrency | One database query per identical concurrent request group |
| Nodes | 177 | 72 (published result) |
| Change | Operate in place | Migrate trillions of messages, no downtime |
R3.2 What Breaks in the Round 2 Design
| Round 2 piece | What breaks |
|---|---|
| The JVM-based store | GC pauses dominated the p99; compaction backlogs cascade; nodes need babysitting. |
| API servers calling the database directly | Each of thousands of identical reads becomes its own database query; nothing can merge them, because they arrive at different API servers. |
| Operations in place | Maintenance on 177 nodes no longer fits its windows, and the store can't be swapped without moving trillions of rows. |
| The migration itself | Trillions of rows must be copied while both stores take live writes, edits and deletes, without losing, duplicating or resurrecting any of them. |
R3.3 New Requirements and API Additions
An internal data-service API. Discord's data services have "roughly one gRPC endpoint per database query and intentionally contain no business logic". The names below are ours:
| RPC | Request | Returns | Routing key |
|---|---|---|---|
GetMessagesBefore | channel_id, before_id, limit | Up to limit rows, newest first | channel_id |
GetMessagesAfter | channel_id, after_id, limit | Up to limit rows, oldest first (a reverse query) | channel_id |
GetMessage | channel_id, message_id | One row | channel_id |
InsertMessage | full row, write timestamp | ack | channel_id |
UpdateMessageColumns | key, changed columns, write timestamp | ack | channel_id |
DeleteMessage | key, write timestamp | ack | channel_id |
Permissions and product rules stay in the API monolith; the data service only runs queries.
Migration controls (our design; the phases are Discord's):
yamlmessages_store: primary: cassandra # cassandra | scylla dual_write: true # write both stores; the primary's result is returned shadow_read_percent: 1 # compare this share of reads across both stores backfill: source: cassandra target: scylla preserve_write_timestamps: true checkpoint: per-token-range config_version: 17
primarydecides which store's answer the client gets. Flipping it is the cutover.dual_writestarts before the backfill, so everything the backfill doesn't copy arrives through live writes.shadow_read_percentis Discord's "small percentage of reads to both databases".config_versionlets every data-service instance report which settings it runs, so we can see when all of them have flipped.
R3.4 Design Evolution: A New Engine, a Shield and a Migration
Step 3.1: Tail Latency From the Database Runtime
The problem: GC pauses and compaction backlogs dominate our p99, and years of tuning haven't fixed them. We need a store that speaks the same query language and data model, so the application barely changes, but behaves predictably at the tail. What would you do?
Primitive: Write-Ahead Log and LSM Trees (ScyllaDB is also an LSM engine: commit log, memtables, SSTables and compaction) · Loop primitive: Write-Ahead Log, fsync & Group Commit (what an acknowledged write has to survive, and why the durable half of the mirror matters)
Step 3.2: Thousands of Identical Reads Hit One Partition
The problem: an @everyone announcement. 5,000 requests a second (assumption) ask for exactly the same thing: the latest 50 messages of one channel. The API monolith runs on many servers, so these requests arrive all over the fleet. Each becomes its own database read. What would you do?
Synthesizing vector architecture diagram...
Three requests, one query. Client 2 waits 8 ms and client 3 waits 5 ms: joiners wait only for the remainder of the query already in flight, never longer than one query.
Primitive: Distributed Cache Patterns and Eviction (coalescing, single-flight and stampedes) · Primitive: Consistent Hashing (routing by channel) · Drill: The Product Page That Melted Redis (answered here: a cold-cache stampede is merged by coalescing concurrent misses; "cache it forever" is the second wrong answer above)
Step 3.3: Migrate Trillions of Messages Without Downtime
The problem: trillions of rows on nearly 200 Cassandra nodes. Live writes, edits and deletes never stop. We need every message in ScyllaDB, none lost, none duplicated, none resurrected, then a switch users don't notice. What would you do?
Primitive: Change Data Capture and the Outbox Pattern · Loop primitive: Change Streams & the Transactional Outbox (Part 1: why writing to two systems can disagree; Part 7: refill from the stored source) · Loop primitive: Leases, Fencing Tokens & Distributed Locks (Part 9: when idempotent writes let one worker win without a lock)
Trace: one message during the migration
Synthesizing vector architecture diagram...
Order doesn't matter: whichever arrives first, each cell ends with its newest version, because both paths carry the source write time.
Notice the moment between the dual-written edit and the backfill: ScyllaDB holds a row for m with content and no author_id. That's exactly the shape Round 2's cleanup treats as "an edit raced a delete; remove it". So that cleanup stays off on ScyllaDB until the backfill is complete and shadow reads are clean, including after the flip. Shadow reads compare; they don't repair. And any delete the cleanup issues on Cassandra during the migration must be dual-written like every other delete, or ScyllaDB keeps the broken row.
Step 3.4: Deletes Still Hurt
The problem: the migration stalled on uncompacted tombstone ranges. ScyllaDB has no GC pauses, but a read that scans a million tombstones is still a read that scans a million tombstones. What would you do?
Loop primitive: LSM-Trees & Compaction (Part 5: compaction merges and when a tombstone may go)
Step 3.5: Can We Archive Old Messages?
The problem: most channels' old buckets are almost never read, yet every byte sits on three replicas on fast local disks. Could we move old history to object storage and load it on demand? What would you do?
Round 3 Step Summary
| Step | Problem | Component | What it costs us |
|---|---|---|---|
| 3.1 | Tail latency from the runtime | ScyllaDB: C++, shard-per-core, no GC; faster reverse queries; the super-disk | A migration; new storage operations |
| 3.2 | Identical concurrent reads | Rust data services, request coalescing, consistent hashing by channel | A new tier and hop |
| 3.3 | Migrate trillions, no downtime | Dual writes, Rust migrator with checkpoints, shadow-read validation, flip | Two clusters for weeks |
| 3.4 | Deletes still hurt | Repair-based tombstone GC, short pages, kinder delete patterns | Repair gates space reclaim |
| 3.5 | Archive old messages? | Evaluated; not built | A second read and delete path, if built |
R3.5 Global Architecture
Synthesizing vector architecture diagram...
Every query now passes through a data service that routes by channel and merges identical concurrent reads; only then does it reach ScyllaDB. The dotted parts are the migration (done) and the archive (an option Discord said it hoped to avoid).
Our AWS translation (not Discord's setup). ScyllaDB on storage-optimized EC2 instances with local NVMe instance storage (for example the i4i family), 24 nodes in each of 3 AZs with one replica per AZ; data services on EKS; S3 for backups; CloudWatch for alarms. The super-disk idea maps to a RAID 1 between a local NVMe stripe and an EBS volume marked write-mostly; that's a translation of Discord's GCP design that we haven't seen Discord or AWS publish, so it would need its own testing. EC2 instance storage has the same property that drove Discord's design: its data doesn't survive the instance being stopped or its host failing.
Trace: an old-message jump
Synthesizing vector architecture diagram...
A jump into history is two range reads in one old bucket, one of them in reverse order: the query Discord needed ScyllaDB to make fast.
Sources for this round
- Ingram, How Discord Stores Trillions of Messages, Discord blog, March 2023: reasons for ScyllaDB; all other databases migrated by 2020; reverse queries; data services, coalescing and consistent-hash routing by channel; the super-disk cluster; dual writes, the Spark estimate, the Rust migrator (SQLite checkpoints, 9 days, 3.2M/s), the tombstone stall; shadow-read validation; the May 2022 cutover; 72 nodes, 9 TB each; p99 results; the World Cup final.
- Oakley, How Discord Supercharges Network Disks for Extreme Low Latency, Discord blog, August 2022: Local SSD and persistent disk trade-offs on GCP; md RAID 0 plus RAID 1 write-mostly; 375 GB local SSDs; bad-sector handling.
- Ingram, How Discord Migrated Trillions of Messages from Cassandra to ScyllaDB, ScyllaDB Summit 2023 talk: the same story, including the hybrid-RAID1 storage topology and the Rust data service library and migrator.
- Vishnevskiy, How Discord Stores Billions of Messages, January 2017: the archive idea "we want to avoid".
- ScyllaDB documentation: Data definition:
tombstone_gc(modestimeout,repair,disabled,immediate), Compaction strategies, Replace a dead node, nodetool toppartitions, and thequery_tombstone_page_limitsetting in ScyllaDB's configuration source (checked September 2026).
Discord hasn't published its data-service fleet size, its compaction strategy, its per-AZ layout, the migrator's read consistency or how it handled backfill timestamps and tombstone GC during the load; everything on this page about those is our design.
R3.6 Numbers and Cost
Figures marked published come from the sources above; everything else is our assumption or derived from one.
Before and after (published)
| Cassandra (early 2022) | ScyllaDB (reported 2023) | Change | |
|---|---|---|---|
| Nodes | 177 | 72 | 105 fewer, −59% |
| Disk per node | ~4 TB average | 9 TB | 2.25× |
| Total disk | 177 × 4 ≈ 708 TB | 72 × 9 = 648 TB | about the same |
| History fetch p99 | 40–125 ms | 15 ms | 2.7× to 8.3× lower |
| Insert p99 | 5–70 ms | 5 ms, "steady" | the spread is gone |
Synthesizing vector architecture diagram...
The biggest win isn't the median but the range: the Cassandra figures were ranges (40–125, 5–70); the ScyllaDB ones are single numbers.
How full is the new cluster? (our estimate) Suppose ScyllaDB holds 150 TB of unique, compressed message data (assumption). Three copies are 450 TB: 450 ÷ 648 = 69% of the disk, leaving room for compaction to work. If the cluster held "trillions" (say 2 trillion) messages in that, that's 150 TB ÷ 2 trillion = 75 bytes per message on disk. Discord didn't publish the fill level or the message count, so treat this as a reminder of how compact the data must be, not a fact.
The migration (published figures, our arithmetic)
| Quantity | Arithmetic | Result |
|---|---|---|
| Spark migrator estimate | Published | ~3 months ≈ 90 days |
| Rust migrator estimate | Published | 9 days = 777,600 s |
| Speed-up | 90 ÷ 9 | 10× |
| Peak rate | Published | 3.2M messages/s |
| Most it could copy in 9 days at the peak rate | 3,200,000 × 777,600 | ≈ 2.5 trillion |
So the 9-day estimate is consistent with a data set in the low trillions; the post doesn't give the exact count.
The load test nobody planned. Discord's 2023 post shows message sends during the December 2022 World Cup final: one spike per goal and per break, with the store "handling it perfectly". Discord didn't publish the peak rate, so we don't use it as a number; it's the published evidence that the coalescing tier and the new cluster held up under a global, synchronized spike.
Coalescing for one announcement (from step 3.2): 5,000 identical requests/s become 98 queries/s when routed to one instance, 283 with per-AZ routing, 2,500 without routing.
Data-service fleet (assumptions: 100,000 fetches/s and 93,000 writes/s at the peak; 20,000 requests/s per instance; survive losing an AZ)
| Arithmetic | Instances | |
|---|---|---|
| Peak requests | 100,000 + 93,000 | 193,000/s |
| Two AZs must carry it all | 193,000 ÷ 2 = 96,500/s per AZ ÷ 20,000 = 4.8, round up | 5 per AZ, 15 in all |
Load on ScyllaDB at the peak (72 nodes, 24 per AZ)
| Arithmetic | Per node | |
|---|---|---|
| Replica writes | 93,000 × 3 = 279,000/s ÷ 72 | 3,875/s |
| Replica reads | 100,000 fetches × 70% left after coalescing (assumption) × 1.3 buckets × 2 replicas = 182,000/s ÷ 72 | 2,528/s |
Cross-AZ traffic in the AWS translation. Data moving between AZs costs $0.01/GB in each direction, so $0.02 for every GB that crosses. Assume 1 KB per replicated write on the wire, 20 KB per 50-message fetch response, and an average of 50,000 fetches/s (half the peak).
| Stream | Arithmetic | A month (30 days) |
|---|---|---|
| Writes to the 2 replicas in other AZs | 4B messages a day (published) × 1 KB × 2 = 8 TB/day × $0.02/GB = $160/day | ≈ $4,800 (unavoidable with one replica per AZ) |
| API → data service, one global ring | 50,000/s × 20 KB = 1 GB/s = 86.4 TB/day; ⅔ crosses an AZ = 57.6 TB × $0.02 = $1,152/day × 30 | ≈ $34,560 (saved by a ring per AZ) |
| Data service → a coordinator in another AZ | Same ⅔ of the same bytes | ≈ $34,560 (saved by an AZ-aware driver) |
| With both | Coordinator is the local replica; the remote replica sends only a digest | ≈ $0 of read bytes |
AZ-aware coordinators save about $34,560 a month and a ring per AZ another $34,560: about $69,000 a month together. The two are independent: a rack- and token-aware driver makes a local replica the coordinator even with one global ring, at no cost; the ring per AZ costs 3 queries instead of 1 per coalesced group. The per-AZ ring is our variant; Discord didn't publish how its routing treats zones.
Why no load balancer in front of the data services. A Network Load Balancer would process the same ~1 GB/s: 3,600 GB an hour, which is 3,600 NLCUs at 1 GB processed per NLCU-hour for TCP, × $0.006 per NLCU-hour (us-east-1) = $21.60 an hour × 720 hours (30 days, the month used across this page) = ≈ $15,550 a month by the bytes dimension alone. Worse, it balances connections, not channels, so it would scatter a channel's requests and undo coalescing. Routing belongs in the client library.
Running two clusters at once (quota check). During the migration both clusters run. In our translation, 72 new nodes of, for example, i4i.8xlarge (32 vCPUs each) need 72 × 32 = 2,304 vCPUs on top of the old cluster. Storage-optimized instances count toward the account's "Running On-Demand Standard instances" vCPU quota, so we check that quota in the region before provisioning (R3.9, command 7). 72 nodes is 24 per AZ, a whole number, so each AZ holds exactly one replica of every partition.
The archive, priced (translation, list prices in us-east-1). Suppose 60% of the 150 TB is cold (assumption): 90 TB.
| Option | Arithmetic | A month |
|---|---|---|
| S3 Glacier Instant Retrieval | 90,000 GB × $0.004 | $360 |
| S3 Standard, tiered | 50,000 GB × $0.023 + 40,000 GB × $0.022 = $1,150 + $880 | $2,030 |
| Keep in ScyllaDB | 90 TB × 3 copies = 270 TB of NVMe, 270 ÷ 9 = 30 nodes' worth of disk | Disk for 30 of 72 nodes |
Glacier Instant Retrieval bills objects smaller than 128 KB as 128 KB, charges a 90-day minimum, and charges $0.03/GB to retrieve plus $0.01 per 1,000 GETs: hence one file per channel per year, not per bucket. The saving is real only if the cluster can actually shrink, which depends on throughput too: freeing 270 TB takes the cluster from 69% to (450 − 270) ÷ 648 = 28% full, but the 72 nodes may still be needed for requests. That's the quantitative version of "we want to avoid doing this one".
R3.7 Trade-Offs
Change the database, or change the access pattern?
| Change the database (ScyllaDB) | Change the access pattern (data services) | |
|---|---|---|
| Fixes | GC pauses, compaction and repair speed, per-core isolation | Unbounded concurrency on hot partitions |
| Doesn't fix | Hot partitions ("can still be a thing") | Runtime pauses and maintenance toil |
| Cost | A migration of trillions of rows | A new tier and hop |
Discord did both, in the cheaper order: the data services first (they helped even on Cassandra), then the database. Either alone would have left half the problem.
A data-service tier vs direct database access
| Direct from the API monolith | Data services (chosen) | |
|---|---|---|
| Identical concurrent reads | Each is a query | One query per group |
| Concurrency on a hot partition | Unbounded | Bounded per channel by the coalescer |
| Where queries live | Scattered through the monolith | One endpoint per query, no business logic |
| Cost | Nothing extra | A fleet (15 instances by our sizing), one hop, a routing library |
Migration strategies
| Time-based cutover (Discord's first plan) | Everything at once (chosen once the migrator took 9 days) | |
|---|---|---|
| Value | New data on ScyllaDB early | All at once, later |
| Complexity | Reads must stitch two stores by time | One flip |
| Risk | Two sources of truth for months | Two clusters for weeks; verification before the flip |
Spark migrator vs a custom one. Discord's case is the interesting one: the generic tool was the safer choice until its estimate (three months) was longer than the firefighting on the old cluster could wait. Owning a fast data library made a custom migrator an afternoon's work. Without that library, the right answer would have been the three months.
Run ScyllaDB ourselves, or a managed service? Discord runs its own. In an AWS translation, the managed Cassandra-compatible option is Amazon Keyspaces. The trade: no nodes, repairs or compaction to run, against less control over storage layout, compaction and tail behavior, and a pricing model per request instead of per node. For a team whose main complaint was toil, it's a question worth asking; for a team whose main complaint was p99 and who had already moved every other cluster to ScyllaDB, staying self-managed was consistent.
What changed from Round 1. Round 1 matched the layout to the query: a partition per channel and bucket, newest first. Round 3 keeps that layout exactly and changes everything around it: the engine underneath (no pauses, a core per shard, a cache of its own) and a tier above that shapes the traffic before it reaches the partition. The lesson of the whole loop: the data model was right in 2016; what failed was the runtime and the lack of any limit on how hard one partition could be hit.
R3.8 Failure Modes
| Trigger | What you'd see | How the design responds |
|---|---|---|
| A coalesced read misses a just-acknowledged write | A client opening a channel doesn't see a message sent milliseconds earlier | Bounded by one query's duration; the gateway pushes the message anyway. Within an instance, a per-channel write counter stops new requests from joining a read that started before a write completed (step 3.2). |
| A hot partition overloads one core | One shard on 3 nodes at 100% CPU; its other partitions slow; the rest of each node is fine | Coalescing bounds identical reads; rate limits bound distinct ones; alarm per shard, not per node; nodetool toppartitions finds the key. |
| A data-service instance dies | The ring moves its channels to other instances; coalescing for those channels restarts cold | A brief rise in database queries for those channels; the per-AZ ring keeps it within one AZ. Scale instances so a loss fits (R3.6). |
| Verification mismatch | Shadow reads disagree | First re-read both stores after a second: a write landing between the two reads is a race, not a bug. Persistent mismatches go to a list; repair each by re-copying that key from the source of truth (Cassandra, before the flip) with original timestamps. |
| A dual write succeeds on the primary and fails on the secondary | Drift that verification may catch only by sampling | Record each failed secondary write (key and timestamp) in a repair queue, and re-copy the key from the primary later. A repair batch that keeps failing is split in halves until the single bad key is found and set aside, so one poison key doesn't block the rest. |
| A token range is re-copied after a crash, or a backup is restored | Keys that were deleted since the copy or the backup reappear in the target | Original timestamps make re-copies harmless only when the target still has the newer tombstone (hence tombstone_gc disabled during the load). After a restore, compare the restored range with the source and delete keys the source no longer has. ScyllaDB checks no epoch, so a paused old cleanup job could delete a row a newer write just made. The simple guard: issue each cleanup delete USING TIMESTAMP equal to the comparison snapshot's time, so any real write after the snapshot wins by last-write-wins. (The alternative is for the data service to reject deletes from a stale job epoch, or to guard them with an LWT on a row that holds the epoch.) |
| A delayed original request lands after its retry | An edit retried after a timeout, and a later edit, both acknowledged; the delayed original arrives last | Every write carries its timestamp from the data service, set once per logical edit and reused on retries, so the delayed original has the older timestamp and loses (Round 2's rule, still needed). |
| A local SSD fails | On GCP the host is migrated and local SSD data is erased (Discord, 2022) | The persistent-disk half of the mirror still holds the data, and RAID 1 can rebuild the local stripe from it: that's how the design is meant to recover, though Discord's post left the detailed edge cases for a later write-up. Bad sectors on a local SSD are read from the mirror instead of failing the query (published). |
| Tombstone ranges block a scan | A read or a migration times out on a token range | ScyllaDB returns short pages; for Cassandra, compact that range first (Discord's fix at 99.9999%). |
R3.9 Runbook and Incident Response
| Signal | Alarm | Severity | First action |
|---|---|---|---|
| History fetch p99 at the data service | > 30 ms for 5 min (2× the reported 15 ms) | P2 | Which channels? Which shards? |
| Insert p99 | > 10 ms for 5 min | P2 | Disk latency on the super-disks? Compaction? |
| Coalescing ratio (requests ÷ database queries) per instance | Falls sharply during a spike | P3 | Is routing by channel working, or did the ring split? |
| CPU per shard (core) | One shard > 90% for 5 min | P2 | toppartitions for the hot key |
| Pending compactions | Rising for 30 min | P3 | Compaction throughput; disk space |
| Tombstones scanned per read | Rising for one table | P3 | Recent mass deletes? Empty-bucket tracking? |
| Repair completion | Not completed on schedule | P2 | Space isn't being reclaimed in repair mode; in timeout mode, resurrection risk |
| RAID state of each super-disk | Any degraded mirror | P2 | Rebuild the local stripe from the durable disk |
Replace a dead node REL 11 · OPS 10
-
Confirm it's down, not slow:
nodetool statusshows itDNfrom several other nodes (command 1). Note its Host ID. -
Decide: return or replace. If the table uses
timeouttombstone GC and the node has been down longer thangc_grace_seconds, never let it rejoin with its data; replace it. Inrepairmode that particular risk is gone, but a node that missed a long stretch of writes is still quicker to replace than to repair. -
Start a new node with the same configuration and this line in its
scylla.yaml, naming the dead node's Host ID (ScyllaDB's documented procedure):yamlreplace_node_first_boot: 675ed9f4-6564-6dbd-ca08-43fddce952de -
Let it stream its data from the other replicas, then repair it (command 6), unless repair-based node operations are enabled for replace, in which case the replace already repaired.
-
Watch p99, pending compactions and hints drain before calling it done.
Go deeper: CLI playbook
Commands an on-call engineer runs one at a time; replace names with real ones.
text# 1. Node states and Host IDs, seen from this node nodetool status # 2. Pending and running compactions on this node nodetool compactionstats # 3. The hottest partitions of messages.messages over 10 seconds nodetool toppartitions messages messages 10000 # 4. Read and write latency percentiles for requests this node coordinated nodetool proxyhistograms # 5. Table statistics, including SSTable counts nodetool tablestats messages.messages # 6. Repair through ScyllaDB Manager sctool repair -c messages-prod # 7. (AWS translation) The Running On-Demand Standard instances vCPU quota, before provisioning a second cluster aws service-quotas get-service-quota --region us-east-1 --service-code ec2 --quota-code L-1216C47A
R3.10 Pillar Check
| Pillar | What Round 3 adds |
|---|---|
| Reliability | Hot partitions contained by coalescing and per-core isolation; a migration with dual writes, checkpoints, shadow-read verification and a rollback window; durable mirrors under local disks; repair-gated tombstone GC REL 10 · REL 8 · REL 9 |
| Performance Efficiency | A store without GC pauses; reads from local NVMe; one query per group of identical reads; routing by channel; fast reverse queries PERF 1 · PERF 3 |
| Security | Permissions checked in the API before any data-service call (the data services hold no business logic) SEC 3; message content treated as the most sensitive user data we hold, so a delete stays a delete through migration, restores and any archive SEC 7 |
| Cost Optimization | 177 → 72 nodes; AZ-aware coordinators and a ring per AZ save about $69,000 a month of transfer together in our translation; no load balancer on the hot path; the archive priced and declined COST 6 · COST 8 · COST 11 |
| Operational Excellence | Per-shard and coalescing signals; a documented node replacement; a quiet on-call ("We're not having weekend-long firefights", Discord) OPS 8 · OPS 10 · OPS 11 |
| Sustainability | 105 fewer nodes for the same data; duplicate reads never reach the disk SUS 5 · SUS 3 |
R3.11 Round 3 Rubric and Follow-Ups
What an architect (L7) answer adds over L6
- Separates what a new database fixes (runtime pauses, repair speed, isolation) from what it doesn't (hot partitions), and fixes both.
- Designs request coalescing with its routing, quantifies the gain, and states its staleness bound.
- Plans a live migration as dual writes, backfill, verification, cutover and rollback, and names the correctness rules: original timestamps, tombstones kept during the load, caught-up source replicas, idempotent re-copies.
- Reads published results critically: ranges vs single numbers, capacity vs usage, what a 9-day estimate implies.
- Prices the options (transfer, load balancers, archives) and knows when not to build something.
- Keeps Discord's published facts separate from its own design.
Follow-up questions
-
"Why Rust for the data services and the migrator?" Answer: Discord's stated reasons: "fast C/C++ speeds without having to sacrifice safety", the Tokio async runtime, driver support for Cassandra and ScyllaDB, and code that is easy to write safely under concurrency, which is exactly what a coalescer is. And because the migrator reused the data-service library, it took "an afternoon".
-
"Discord runs this in one place. What if we had to run it in two regions?" Answer: Discord hasn't published a multi-region layout for messages, so this is our design. ScyllaDB would replicate asynchronously between data centers, with
LOCAL_QUORUMin each, so the RPO on losing a region is the replication delay, usually seconds, not zero. Deciding a region is lost must not depend on that region, so detection runs from the other region and from outside. Moving users (DNS or routing) should be a gated, human-approved step, checked against capacity in the survivor and any data-residency rules. Any conditional write that must be globally consistent paysSERIALPaxos across regions;LOCAL_SERIALis only consistent within one. For RPO 0, writes would have to wait for copies in both regions (EACH_QUORUMin Cassandra and ScyllaDB), adding a cross-region round trip to every send; the chat loop's Round 3 designs that synchronous alternative and prices it. -
"Isn't request coalescing just a cache with a very short TTL?" Answer: close, but with one important difference: it never serves a result to a request that arrives after the query finished. So there's no expiry to tune, no invalidation, and no stale data beyond one query's flight time.
-
"The p99 went from 40–125 ms to 15 ms. How much was the database and how much the data services?" Answer: Discord didn't split it, and the data services arrived first, while still on Cassandra, where they helped but "don't solve all of our problems". An honest answer says the published numbers can't separate the two.
Interview gotchas from this round's wrong answers
| Gotcha | Why it's wrong |
|---|---|
| "A new database fixes hot partitions" | Discord: they "can still be a thing in ScyllaDB". |
| "Cache in every API server" | Fills per server, invalidates per server; still many queries. |
| "Coalesce without routing" | Identical requests land on different instances; the gain nearly vanishes. |
| "Stop the world to migrate" | Days of downtime at millions of rows a second. |
| "Backfill with today's timestamps" | Old copies beat newer edits and bring deleted messages back. |
| "Archive everything old" | A second read path and a second delete path, for a saving that may not shrink the cluster. |
Loop Closer: Interview Strategy for All Three Rounds
How to Run Each 60-Minute Round
| Time | Round 1 | Round 2 | Round 3 |
|---|---|---|---|
| 0–5 min | Scoping: what broke on MongoDB, the main read, history forever | Restate Round 1 in 60 seconds | Restate Round 2 in 60 seconds |
| 5–15 min | Requirements and the history API with ID cursors | Scope raise → what breaks | Scope raise → what breaks |
| 15–40 min | Steps 1.0–1.4: one replica set → partition by channel → Snowflakes → 10-day buckets → newest first | Steps 2.1–2.5: hot partitions → tombstones → GC and compaction → LWW races → toil | Steps 3.1–3.5: ScyllaDB → data services and coalescing → migration → deletes → archive |
| 40–50 min | Write and read rates, bucket sizing, the storage sanity check | Little's law on a hot partition, tombstone and GC arithmetic | Before/after, migration throughput, coalescing math, fleet and transfer costs |
| 50–60 min | Failures and pillar check | Failures, gotchas, pillar check | Failures, runbook, pillar check |
For how to spend a single 45-minute round, see the 45-minute interview blueprint. For delivering messages to clients (gateways, fan-out, huge channels), see the chat loop; for quorums, hinted handoff and anti-entropy from the inside, the key-value store loop; for Snowflake clocks in depth, the unique ID loop.
The Two Sentences That Matter Most
- Opening any round: "Every read is 'the latest messages in one channel', so I'll partition by channel and a time bucket, sort by a time-ordered ID newest first, and size the bucket so the busiest channel stays small."
- When scale arrives: "A new engine fixes pauses but not hot partitions, so I'll put a data service in front that routes by channel and turns identical concurrent reads into one query, then migrate with dual writes, a timestamp-preserving backfill and shadow-read verification."
Well-Architected Review Sheet
Interviewers rarely ask "which pillar is this?". They ask the pillar's question in plain words. Rehearse one sentence per row.
| Pillar | Question you'll hear | One-sentence answer | Round | Backed by |
|---|---|---|---|---|
| Reliability | "What if a node dies?" (REL 11) | Three copies, quorum reads and writes, hints for 3 hours and repair after that; a node down past the grace period is rebuilt. | 1–2 | R1.9, step 2.2 |
| "One channel is melting the cluster. Why, and what do you do?" (REL 10) | One key lives on three replicas; bound its concurrency with coalescing routed by channel, and isolate it per core. | 2–3 | Steps 2.1, 3.1, 3.2 | |
| "How do you change the database under live traffic?" (REL 8) | Dual writes, a checkpointed backfill with original timestamps, shadow reads compared, then a flip with a rollback window. | 1, 3 | R1.6, step 3.3 | |
| Performance | "How do you make 'latest 50' fast?" (PERF 3) | Partition by channel and 10-day bucket, cluster by Snowflake descending: one partition, adjacent rows. | 1 | Steps 1.1–1.4 |
| "Why did tail latency improve?" (PERF 1) | No GC pauses, a shard per core, local-disk reads, and one query per group of identical reads. | 3 | Steps 3.1, 3.2 | |
| Security | "Messages are sensitive user data. When a user deletes one, is it gone?" (SEC 7) | After the grace period and a compaction, yes; and the design stops repairs, migrations and restores from bringing it back. | 2–3 | Steps 2.2, 3.3, R3.8 |
| Cost | "What did the old system really cost?" (COST 11) | Engineers' time: pages, gossip dances, reboots, repairs cut back. | 2 | Step 2.5 |
| "Where do transfer charges hide?" (COST 8) | In routing across zones: AZ-aware coordinators (≈$34,560 a month) and a ring per AZ (another ≈$34,560) save about $69,000 a month together in our translation. | 3 | R3.6 | |
| Operations | "How do you know it's healthy?" (OPS 8) | p99 per query at the data service, CPU per shard, coalescing ratio, pending compactions, tombstones per read, repair completion. | 2–3 | R2.8, R3.9 |
| Sustainability | "How did you do more with less?" (SUS 5) | 177 nodes became 72 for the same data, and duplicate reads never touch a disk. | 3 | R3.6 |
Rubric Across Levels
| Dimension | L5 (Round 1) | L6 (Round 2) | L7 (Round 3) |
|---|---|---|---|
| Data model | Partition by channel and bucket; Snowflake clustering, newest first | Knows what the model costs: hot keys, tombstones, wide partitions | Keeps the model; changes the engine and the traffic around it |
| Consistency | Quorum reads and writes | LWW per column; LWT and SERIAL vs LOCAL_SERIAL; client-side timestamps | Timestamp-preserving backfills; staleness bound of coalescing |
| Deletes | Deletes exist | Tombstones, grace period, resurrection, empty buckets | Repair-based tombstone GC; tombstones through a migration and a restore |
| Operations | Dark launch | Measures toil; gossip dance; node down past grace | Node replacement, per-shard signals, a quiet on-call |
| Honesty | Labels assumptions | Separates Discord's measures from its own mitigations | Reads published results critically; says what the numbers can't show |
Sources
All sources used on this page, oldest first.
- Vishnevskiy, How Discord Stores Billions of Messages, Discord blog, January 2017.
- Oakley, How Discord Supercharges Network Disks for Extreme Low Latency, Discord blog, August 2022.
- Ingram, How Discord Stores Trillions of Messages, Discord blog, March 2023.
- Ingram, How Discord Migrated Trillions of Messages from Cassandra to ScyllaDB, ScyllaDB Summit 2023.
- Discord Developer Documentation, API Reference: Snowflakes and Message resource, checked September 2026.
- Apache Cassandra documentation, cassandra.apache.org/doc: tombstones,
gc_grace_seconds, hinted handoff, compaction, lightweight transactions. - ScyllaDB documentation: Data definition (
tombstone_gc), Compaction strategies, Replace a dead node, nodetool toppartitions, and the configuration source forquery_tombstone_page_limitandmax_hint_window_in_ms, checked September 2026. - AWS, Amazon S3 pricing (
us-east-1, checked September 2026), and the EC2 data transfer and Elastic Load Balancing pricing pages, for the translation's cost figures.