KnowSys

Designing Discord

Aiko posts one line in #general of a game's official server while three hundred thousand other people have it open. We'll design the system that puts her message on all their screens within a moment and keeps it, among trillions of others, for as long as Discord exists: from one database to bucketed partitions, Snowflake IDs, tombstones, ScyllaDB and request coalescing, then Elixir guild processes, relays, passive sessions, lazy member lists, compressed WebSockets and read states.

⏱ 50 min read◆ IntermediateAssumes: the WhatsApp case study (chapter 49), chapter 29 (partitioning) helps, chapter 18 (storage engines) helps
Start reading

Aiko is at her desk in Osaka with Discord open. On the left of the window is a column of round icons, one for each community she belongs to, and she has the official server for Starfall, an online game, selected. It has about 2.4 million members. Starfall's game servers have been down for maintenance all afternoon, and the #general channel is full of people waiting. A moderator posts that the servers are back, with an @everyone mention, which notifies every member of the guild. Aiko types "finally!! see everyone in game" and presses Enter. Her line appears at the bottom of the channel, and within a second it scrolls past on the screens of everyone else who has #general open, among a flood of other people saying the same thing.

A huge dark hall packed with rows of glowing computer monitors and people sitting at them
DreamHack Winter 2004 in Sweden, one of the biggest LAN parties of its time: a few thousand players in one hall. A large Discord server is that hall a thousand times over, with hundreds of thousands of people in the same room at once, and every one of them expects to hear what anyone says.Photo: Toffelginkgo, CC BY-SA 3.0, via Wikimedia Commons

In a two-person chat, sending a message means getting it to one other phone. Chapter 49's WhatsApp case study is about doing that well. Aiko's message has a different shape of problem. About 320,000 people are online in the Starfall server at that moment, and Discord's servers need to tell every one of their apps that something happened in #general, while thousands of other people are posting at the same instant. Then the message has to be stored, permanently, because a member who joins next year can scroll back and read it. And each of those 2.4 million members has a little badge counting the messages they haven't read and the ones that mention them, and those have to change too.

So here's the question this case study answers: when Aiko presses Enter, how does Discord get her message onto hundreds of thousands of screens in under a second, and keep it findable among trillions of other messages forever? We'll start with the design Discord itself launched with in 2015, see exactly where it broke, and follow the fixes Discord published over the next decade: how messages are laid out on disk, the database migrations that followed, the processes that fan a message out, and the smaller pieces around them that turned out to matter just as much.

01What we're building, and how big

1.1What it has to do

A few words first, because Discord's vocabulary is a little unusual. What the app calls a server (Starfall's community), Discord's engineers call a guild, and we'll use that word for it so it can't be confused with the computers that run everything. A guild contains channels such as #general, and a channel contains messages. Each member can see some channels and not others, depending on the roles they have in the guild.

Stripped down, the part of Discord we're designing has to do five things:

  1. Deliver a message in real time to every member who's online and allowed to see the channel.
  2. Keep every message forever, and let anyone scroll back through a channel's history, or jump to a message from years ago.
  3. Support edits, deletes and reactions on stored messages.
  4. Show the member list: who's in the guild, grouped by role, with who's online.
  5. Track what each person has read, with an unread marker and a count of mentions per channel.

And the qualities it needs:

  • Fast: a message appears for everyone within a second or so.
  • Huge rooms: a single guild can have millions of members and a million online at once.
  • Permanent: messages are never thrown away, unless someone deletes them.
  • Predictable: loading a channel takes about the same time whether the channel is busy or quiet, new or old.

Two of these break with WhatsApp's design straight away. WhatsApp deletes a message from its servers once it's been delivered (chapter 49, section 3.2), so its storage only ever holds what's in transit. Discord keeps everything, so its storage only grows. And a WhatsApp group tops out at 1,024 members, small enough to copy each message into every member's mailbox. Copying Aiko's message into 2.4 million mailboxes would turn one message into 2.4 million writes, and #general gets several messages a second. So Discord stores one copy of each message, in the channel, and every reader fetches it from there. Chapter 49 calls this fan-out on read, and the rest of this case study follows from it.

1.2How big is it?

Discord says it had more than 90 million people using it every day at the end of 2025. Its engineering posts give the numbers that shape the design. In January 2017 it was storing more than 120 million new messages a day, up from 40 million a day in July 2016. By the start of 2022, its message database held trillions of messages on 177 machines. In October 2023, the largest guild, Midjourney's, had more than 10 million members, over a million of them online at any time, and Discord said that across a few years it had pushed the largest guilds from tens of thousands of online users to almost two million.

Your turn: design it before reading on

Starfall has 320,000 members online, and #general is getting 50 messages a second right after the maintenance ends. If every online member's app must hear about every message, how many deliveries a second is that? Now suppose that on a busy evening every one of those 320,000 people posts one message somewhere in the guild. How many notifications does that make?

That square is the reason big guilds are hard. Storage grows with the number of messages, and that's manageable: add disks. Fan-out grows with messages times listeners, and in a guild both grow together. Sections 5 and 6 come back to it. First, though, the messages have to be stored somewhere, and that's where Discord's first design ran into trouble.

02Version 1: one database with an index

2.1Discord's 2015 design

Discord built its first version in just under two months, in early 2015. Clients connected to a gateway: a fleet of servers that hold one long-lived WebSocket connection per client and push events down it, for the reasons chapter 49 gives in its section 2 (a server that knows first should tell the phone, not wait to be asked). An API server handled requests such as "send this message" and "give me the last 50 messages in #general". Every message went into a single MongoDB database, replicated for safety but not split across machines, in one collection with an index on channel ID and creation time.

Version 1: one database holds every message
sendinsertpublishpushAiko's appOther members320,000 onlineAPI serversGatewayone WebSocket per clientMongoDBindex (channel_id, created_at)
Step 1. Aiko presses Enter. Her app sends the message to an API server.
1 / 4

An index here is a sorted structure, a B-tree, that maps (channel, time) to where each message sits on disk. "The last 50 messages in #general" becomes a walk along 50 neighbouring entries of the index, followed by 50 lookups of the messages themselves. When the index and the messages are in memory, that's fast. In November 2015, at about 100 million stored messages, they stopped fitting in memory, and response times became unpredictable.

?Why does falling out of memory hurt so much?

Because of how people use Discord. Discord's 2017 post on this migration describes three kinds of guild. Voice-chat guilds send a message or two every few days; when someone opens one of their channels, the last 50 messages are scattered across months of the database, and each could be a separate trip to disk. Private guilds of friends send something like 100,000 to a million messages a year but rarely scroll back, so their older messages are almost never in memory. Big public guilds like Starfall send millions of messages a year, but people mostly read the last hour, which does stay in memory. So reads were spread almost randomly over the whole data set, and reads and writes were roughly half and half. Once the data was larger than memory, a large share of reads went to disk, and a disk read costs about a hundred times more than a memory read even on an SSD (chapter 9). Adding memory to one machine only delays the day the data outgrows it again.

2.2What the replacement must do

Discord wrote down what it wanted from the replacement. It had to scale linearly, so that doubling the machines doubles the capacity. It had to fail over automatically when a machine died. It had to need little care, because, in 2017, Discord's backend was four engineers with no dedicated operations staff. It had to be proven technology, and its performance had to be predictable: Discord's alerts fired when the API's 95th-percentile response time, the time within which 95 of every 100 requests finish, went over 80 milliseconds. And Discord didn't want a cache like Redis or Memcached in front of it, since that would be one more system to run and keep consistent.

"Scale linearly" means the data must be split across many machines, each holding a part. Chapter 29 covers the general problem. For us, the question is which part each machine gets, and that choice turns out to decide almost everything about how well the store works.

03Storing messages in Cassandra

3.1Partitions: keep a channel together

Around the turn of 2016 Discord chose Apache Cassandra, a database that spreads data over a ring of machines, as the only one it found that met all of those requirements. Cassandra's data model is the thing to understand first, because Discord's design is mostly a choice of keys.

Every table in Cassandra has a partition key. Cassandra hashes it to a number, called a token, and the token decides which machines store the row. Rows with the same partition key always live together, on the same machines, so reading many of them is one request to one place. Within a partition, rows are kept sorted by a second part of the key, the clustering key, so reading a range of them in order is a single sequential scan.

A ring of eight nodes n1 to n8, with coloured arcs showing that the key foo is stored on n2, n3 and n4 and the key bar on n3, n4 and n5
How Cassandra places data. Each node owns a position (token) on a ring. A key is hashed onto the ring, and its partition is stored on the next few nodes clockwise: here hash(foo) lands on n2, n3 and n4, three copies because the replication factor is 3. Discord's message cluster in 2017 used exactly that factor, so every partition lived on three machines.Image: Apache Cassandra documentation, The Apache Software Foundation, Apache License 2.0

Every question Discord asks of messages is about one channel: the newest 50 in #general, the 50 before this one, the messages around a link from a search result. So the channel ID is the natural partition key, and every channel's messages sit together on three machines. For the clustering key, Discord wanted the messages sorted by time, and that needs an ID that sorts by time. That's where Discord's IDs come in.

3.2Snowflake IDs: a timestamp you can sort

zoomDiscordMessage storePartition keySnowflake ID

Every ID in Discord, for a user, a guild, a channel or a message, is a 64-bit number called a Snowflake, a scheme Twitter published in 2010. Its idea is to build the ID out of the time it was made, so that sorting IDs sorts by time, and to make each one unique without any server asking any other server.

A row of 63 bit cells: 41 marked t for timestamp, 10 marked i for instance, and 12 marked s for sequence
Twitter's original Snowflake layout: a millisecond timestamp in the high bits, then the ID of the machine that made it, then a counter. Discord uses the same shape with 42 timestamp bits (it uses the top bit too), and splits the 10 machine bits into a 5-bit worker ID and a 5-bit process ID.Image: Sasmito Adibowo, CC BY-SA 3.0, via Wikimedia Commons

Discord's developer documentation gives the layout:

BitsFieldWhat it holds
63 to 22 (42 bits)timestampmilliseconds since the Discord epoch, the first moment of 2015 (Unix time 1420070400000 ms)
21 to 17 (5 bits)worker IDwhich machine made the ID
16 to 12 (5 bits)process IDwhich process on that machine
11 to 0 (12 bits)incrementa counter that goes up for every ID that process makes

Three properties fall out of the layout. Because the timestamp is in the high bits, a bigger ID was made later, so sorting by ID sorts by time to the millisecond. Because each process has its own worker and process number, two processes can never make the same ID, so no coordination is needed. And because the counter has 12 bits, one process can make 4,096 IDs in the same millisecond before it has to wait for the next. 42 bits of milliseconds run out about 139 years after 2015.

With that ID, Discord's 2017 table looked like this, simplified in the post to a few of its 16 columns:

SQL
CREATE TABLE messages (
  channel_id bigint,
  message_id bigint,
  author_id bigint,
  content text,
  PRIMARY KEY (channel_id, message_id)
) WITH CLUSTERING ORDER BY (message_id DESC);

Read the primary key (channel_id, message_id) as: partition by channel, and within the channel keep messages sorted by Snowflake, newest first (DESC). "The latest 50 in #general" is now the first 50 rows of one partition.

3.3Partitions that never stop growing

That table has a problem you can spot by thinking about #general a year later. A partition holds every message the channel has ever had, and Discord never deletes anything. When Discord first imported its existing messages, some partitions were already over 100 MB. Cassandra will technically accept partitions up to 2 GB, but a huge partition can't be split across machines, so the three machines holding it carry all of its load. It also makes every background task that touches it, rewriting its files or repairing its copies, slow and memory-hungry. Discord's 2017 post puts it this way: just because it can be done doesn't mean it should.

So cut each channel's history into time slices. Discord looked at its largest channels and found that ten days of messages, even in the busiest, would comfortably stay under 100 MB. So it added a bucket to the partition key: a number for each ten-day window since the epoch, worked out from the Snowflake.

SQL
CREATE 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);

Those double parentheses make (channel_id, bucket) the partition key together. Each channel is now a series of partitions, one per ten days, and each lands on its own three machines somewhere on the ring. Because the bucket comes from the Snowflake's timestamp, anyone holding a message ID can work out which partition it lives in without asking anything.

Two table diagrams: available_rooms_by_hotel_date keyed by hotel_id, and a bucketed version keyed by hotel_id and month
The same technique in Cassandra's own data-modelling guide. K marks partition-key columns and C clustering columns. On the right, adding month to the partition key splits one hotel's ever-growing partition into one per month. Discord's (channel_id, bucket) is the same move with ten-day buckets.Image: Apache Cassandra documentation, The Apache Software Foundation, Apache License 2.0

Reading gets one extra step. To load #general, the API works out the bucket for "now" and reads the newest messages from that partition. If it gets fewer than 50, it moves to the previous bucket, and so on back towards the bucket in which the channel was created (the channel's own Snowflake says when that was). For Starfall's #general, one bucket is always enough. For a quiet channel it isn't.

This program computes buckets the way Discord's 2017 post does: shift the Snowflake right by 22 bits to get its milliseconds, and divide by the number of milliseconds in ten days. It unpacks the example ID from Discord's documentation, buckets a few dates, then looks at three channels with very different traffic, assuming maybe 200 bytes stored per message (an illustration; Discord doesn't publish an average).

Snowflakes, ten-day buckets, and how big a bucket gets
python
Python
import math
from datetime import datetime, timezone
 
DISCORD_EPOCH = 1420070400000            # 1 Jan 2015, in milliseconds
BUCKET_SIZE = 1000 * 60 * 60 * 24 * 10   # ten days, in milliseconds
 
def unpack(snowflake):
    ms = snowflake >> 22                 # top 42 bits: ms since the epoch
    worker = (snowflake >> 17) & 0x1F    # 5 bits
    process = (snowflake >> 12) & 0x1F   # 5 bits
    increment = snowflake & 0xFFF        # 12 bits
    return ms, worker, process, increment
 
def make_bucket(snowflake):
    return (snowflake >> 22) // BUCKET_SIZE
 
# The example ID from Discord's API documentation
sf = 175928847299117063
ms, w, p, inc = unpack(sf)
when = datetime.fromtimestamp((ms + DISCORD_EPOCH) / 1000, timezone.utc)
print(f"id {sf}: {when:%Y-%m-%d %H:%M:%S} UTC, worker {w}, process {p}, increment {inc}")
print(f"  bucket {make_bucket(sf)}  -> partition key (channel_id, {make_bucket(sf)})")
 
# Make an ID for any moment, then bucket it
def snowflake_at(dt, inc=0):
    ms = int(dt.timestamp() * 1000) - DISCORD_EPOCH
    return (ms << 22) | inc
 
for d in ["2015-01-05", "2015-01-11", "2023-03-06", "2026-10-10"]:
    s = snowflake_at(datetime.fromisoformat(d + "T12:00:00+00:00"))
    print(f"{d}: bucket {make_bucket(s)}")
 
# Three channels, about 200 bytes stored per message
for name, per_day in [("#general (giant server)", 300 * 60 * 24),
                      ("#book-club (small server)", 3 * 60 * 24),
                      ("#lfg (voice-heavy server)", 1 / 3)]:
    per_bucket = per_day * 10
    print(f"{name:26} {per_bucket:>11,.0f} msgs/bucket {per_bucket * 200 / 1e6:>8.1f} MB"
          f"   buckets read for the latest 50: {math.ceil(50 / per_bucket)}")
output
C++
id 175928847299117063: 2016-04-30 11:18:25 UTC, worker 1, process 0, increment 7
  bucket 48  -> partition key (channel_id, 48)
2015-01-05: bucket 0
2015-01-11: bucket 1
2023-03-06: bucket 298
2026-10-10: bucket 430
#general (giant server)      4,320,000 msgs/bucket    864.0 MB   buckets read for the latest 50: 1
#book-club (small server)       43,200 msgs/bucket      8.6 MB   buckets read for the latest 50: 1
#lfg (voice-heavy server)            3 msgs/bucket      0.0 MB   buckets read for the latest 50: 16

In the first two lines, the Snowflake is taken apart: the documentation's example ID was made on 30 April 2016 by worker 1, process 0, as that process's eighth ID in that millisecond (the counter starts at 0). Next, the dates show buckets ticking over every ten days from the start of 2015: by today, channels are on bucket 430.

Now the last three lines, where the bucket size cuts both ways. The quiet #lfg channel, with a message every three days, has to read 16 partitions, each a separate request that might go to a different machine, to show 50 messages. And #general at 300 messages a minute (five a second) packs over four million messages, about 860 MB, into a single ten-day bucket, far past the 100 MB that ten days was chosen to stay under in 2017. Ten days was right for the busiest channels of 2017. Guilds kept getting bigger, and the bucket size was fixed in the key.

3.4Hot partitions

Size isn't the only problem with #general's current bucket. For the next ten days every new message in the channel goes to that one partition, and so does every read of the latest messages, almost every read there is. Those requests all land on the same three machines out of hundreds, however many machines the cluster has. A partition that gets far more traffic than its share is called a hot partition, and it's the price of keeping a channel together.

A log-log plot of daily page views against popularity rank for Wikipedia articles, falling roughly in a straight line from millions to one
Popularity is never even. Daily views of Wikipedia articles fall off steeply with rank: the top handful get millions of views, the millionth gets a few. Discord channels behave the same way, so a few channels, and within them the current bucket, take a large share of all traffic.Image: West.andrew.g, CC BY-SA 3.0, via Wikimedia Commons
Decision

What should the partition key for messages be?

message_id alone
Hash every message to its own place on the ring.
  • Load spreads perfectly, even for the busiest channel
  • Reading the last 50 messages of a channel touches up to 50 machines
  • No way to scan a channel in order
channel_id
One partition per channel, sorted by Snowflake.
  • The latest 50 is one sequential read on one partition
  • Partitions grow forever
  • A busy channel's partition is hot for its whole life
chosen
(channel_id, bucket)
One partition per channel per ten days.
  • Reads stay sequential within a bucket
  • Partition size is bounded by traffic per ten days
  • Old buckets go cold and stay cold
  • Quiet channels read many buckets
  • The current bucket of a busy channel is still hot and can still grow large

Discord chose to optimise for the common read, the latest messages in one channel, and to accept hot partitions as a cost to be managed elsewhere. As section 4 shows, that cost grew with the guilds, and the fixes Discord eventually made were in the layers around the database, not in the key. For quiet channels, Discord added a cheaper fix in 2017: remember which of a channel's buckets are empty and skip them.

3.5How Cassandra writes, and what a delete is

To follow what went wrong next, we need to know how Cassandra stores rows on each machine. Like most databases built for heavy writes, it uses a log-structured merge tree (LSM tree), which chapter 18 covers in depth. In short: a write goes into a sorted table in memory, the memtable, and to a log for safety. When the memtable fills, it's written out in one go as an immutable sorted file called an SSTable. Files are never edited. A newer version of a row goes into a newer file, and a read checks the memtable and the files, newest first, merging what it finds.

Sorted files at level 0 merging, with arrows, into fewer, larger sorted files at level 1 and then level 2
Compaction in an LSM tree. New sorted files appear at the top; in the background they're merged into fewer, larger sorted files below. Merging is when old versions of a row are thrown away, and the only time deleted data really leaves the disk.Image: Ben Stopford, CC BY-SA 4.0, via Wikimedia Commons

Left alone, files pile up and every read has to check more of them. So in the background Cassandra merges files together, keeping only the newest version of each row. That merging is called compaction, and it costs disk bandwidth and CPU that would otherwise serve requests.

Now, how do you delete a row from a file you're not allowed to edit? You write a marker into a newer file that says "this row is deleted as of time T". That marker is a tombstone. Reads see the tombstone and hide the row. Compaction eventually drops both the row and the tombstone, but not straight away: a tombstone must survive long enough to reach all three copies of the partition, or a machine that missed the delete could later hand the old row back to the others and bring it back to life. Cassandra keeps tombstones for a grace period, ten days by default, during which a nightly repair compares the copies and fixes any that differ.

Two details of Cassandra's model caught Discord out while it ran both databases side by side. First, every write is an upsert: writing a row that doesn't exist creates it, and the latest write of each column wins. A user editing a message at the same moment another user deleted it could leave behind a row that held only the key and the new text, with no author. Discord's fix was to treat a message with no author as deleted. Second, writing a null into a column is itself a delete, so it writes a tombstone. Discord's code wrote all 16 columns every time, and an average message used only 4, so each message created about 12 needless tombstones. Discord fixed it by writing only the columns that have values.

3.6One message behind millions of tombstones

About six months after the switch to Cassandra, the cluster became unresponsive. Machines were spending ten seconds at a time paused in garbage collection, over and over.

Cassandra runs on the Java virtual machine, and Java frees memory with a garbage collector: from time to time it finds objects the program can no longer reach and reclaims their memory. Some of that work needs the program to stop while it happens, a pause that's normally a few milliseconds.

An animation of a graph of objects in which reachable objects are marked from the roots and the unmarked ones are then swept away
Mark and sweep, the basic garbage-collection algorithm: follow pointers from the roots and mark everything reachable, then free everything unmarked. Real collectors are far more refined, but the cost still grows with how much garbage a program makes, and a program that makes it faster than the collector frees it ends up paused for long stretches.Animation: M17, CC0, via Wikimedia Commons

It was caused by one channel in a public guild for the game Puzzles & Dragons. Someone had used the API to delete millions of its messages, leaving one. Loading that channel took 20 seconds. Can you see why, from the last section?

Predict before you read on

The channel holds one live message and millions of recently deleted ones. Why does reading the latest 50 messages take 20 seconds and make the machine pause?

Discord made two fixes. It cut the tombstone grace period from ten days to two, which was safe because repairs ran every night, so a delete reached all copies well within two days. And it made the read path remember which buckets of a channel are empty and skip them, so a channel emptied by mass deletion costs at most one bucket's scan. Discord also noted the deeper worry: big partitions put pressure on the garbage collector during compaction too, and nothing about that would get better as channels grew.

By January 2017 Discord had 12 Cassandra machines, each copy of the data stored three times, with nearly a terabyte of compressed data per machine. Reads took under 5 milliseconds and writes under one in testing, and Discord moved the rest of its live data onto Cassandra. What nobody knew yet was how that would hold up as "billions" became "trillions".

04Trillions of messages: ScyllaDB and the data services

4.1What 177 machines felt like

By early 2022 the message cluster had 177 Cassandra machines holding trillions of messages, and Discord's 2023 post describes it as a "high-toil" system: the on-call engineers were paged often, and latency was unpredictable. Every problem we've met so far had grown with it.

Hot partitions spread their pain. Discord reads and writes at quorum: two of the three copies must answer before a request succeeds (chapter 28). When a big guild's current bucket got hammered, the three machines holding it slowed down, and every other request that needed those machines as part of its quorum slowed down too, including requests for unrelated quiet channels. Compaction fell behind on busy machines. That left more files for every read to check, so reads got slower, so less capacity was left for compaction. Discord's way out was what it called a "gossip dance": take a machine out of rotation, let it compact without serving traffic, put it back, repeat. And garbage-collection pauses were long enough that an engineer sometimes had to reboot a machine by hand.

None of these had a single fix inside Cassandra. Discord's answer came in two parts: a different database underneath, and a new layer in front of it.

4.2ScyllaDB: the same model, no garbage collector

ScyllaDB is a reimplementation of Cassandra in C++. It speaks the same query language and stores data the same way, as LSM trees on a ring, so Discord's tables and keys could move over unchanged. Three things about it appealed to Discord. It has no garbage collector, because C++ programs free memory explicitly. It splits each machine's data by CPU core, a design ScyllaDB calls shard-per-core: each core owns its own slice of the data and its own memory and never shares them, so one hot partition can only overload one core instead of the whole machine. And it promised faster repairs, the nightly job that had grown heavier with every terabyte.

Discord didn't start with messages. By 2020 every other database it ran had moved to ScyllaDB, and messages stayed on Cassandra, because Discord's tests found ScyllaDB too slow at reading a partition in reverse order. Reading newest first is almost every message query, so that was a blocker until the ScyllaDB team improved it.

Discord also changed what sits under the database. Cloud machines offer local SSDs, which are fast but lose their contents if the machine dies, and network-attached disks, which are durable but slower. Discord combined them in what it called a super-disk: reads are served from the local SSDs, and every write is mirrored with software RAID onto a network disk, so the data survives a lost machine.

4.3Data services and request coalescing

Half of the fix doesn't involve the database at all. Go back to Aiko's moment: the moderator posts that the game is back, mentioning everyone, and in the next few seconds tens of thousands of people click on #general. Every one of their apps asks the API for the latest 50 messages in the same channel. Each request turns into a read of the same partition on the same three machines. Tens of thousands of identical questions, each answered from scratch.

Discord's fix, added between its API and its databases, is a set of small services it calls data services, written in Rust. Each one offers roughly one endpoint per database query, called over gRPC (a protocol for calling a function on another machine, chapter 63), and contains no business logic. What makes them useful is two features.

The first is request coalescing. If a data service is already waiting on the database for "latest messages in #general" and a second request for the same thing arrives, it doesn't send a second query. It attaches the second request to the one already in flight, and when the answer arrives, it goes to everyone who asked. However many people ask in those few milliseconds, the database sees one query.

A second feature makes the first one work. Coalescing only helps if identical requests reach the same data service instance. So each request carries a routing key, for messages the channel ID, and the API picks the instance by consistent hashing on that key (chapter 29). Every request for #general goes to the same instance, the one place where they can meet and be merged.

Reading #general through a data service
RUST DATA SERVICES1 queryMembers' appsopening #generalAPImonolith×NData service Achannels hashed hereData service Bowns #generalData service CScyllaDB(channel_id, bucket)
Step 1. Thousands of members open #general at once. Each app asks an API server for the latest 50 messages.
1 / 4

How much does coalescing save? It depends on how many requests arrive while one query is in flight. This program simulates readers opening #general after an @everyone mention, arriving at random over the next few seconds with most in the first half-second, and a database query that takes 5 ms. It counts how many queries reach the database with and without coalescing.

Request coalescing when everyone opens the same channel
python
Python
import random
random.seed(7)
 
QUERY_MS = 5          # one database read of the latest messages takes 5 ms
 
def plain(arrivals):
    return len(arrivals)                  # every reader runs its own query
 
def coalesced(arrivals):
    queries, busy_until = 0, -1.0
    for t in sorted(arrivals):
        if t >= busy_until:               # nothing in flight: start a query
            queries += 1
            busy_until = t + QUERY_MS
        # otherwise: wait for the query already in flight and share its answer
    return queries
 
# Someone posts @everyone in #general. Readers open the channel over the next
# few seconds; most arrive in the first half-second.
for readers in [100, 10_000, 200_000]:
    arrivals = [random.expovariate(1 / 500) for _ in range(readers)]  # ms
    print(f"{readers:>7,} readers: {plain(arrivals):>7,} queries without coalescing, "
          f"{coalesced(arrivals):>4,} with")
output
C++
    100 readers:     100 queries without coalescing,   70 with
 10,000 readers:  10,000 queries without coalescing,  461 with
200,000 readers: 200,000 queries without coalescing,  762 with

Look at how the "with" column grows. With 100 readers spread over a few seconds, most arrive when nothing is in flight, so coalescing saves little: 70 queries. With 10,000 readers, it saves 95%. With 200,000, it saves more than 99.6%, and the number of queries hardly grows at all, because it can never exceed the length of the burst divided by the query time, a few hundred 5 ms windows. Coalescing does the most for the hottest keys, exactly the ones that hurt. The database stops seeing how popular a channel is.

4.4Moving trillions of messages

With the data services in place, Discord had to move the data itself, without downtime. It started writing every new message to both Cassandra and ScyllaDB, so that from a cut-over point onwards ScyllaDB had everything new. Old messages still had to be copied. ScyllaDB's own migration tool, built on Spark, estimated three months. Discord's engineers wrote a replacement in Rust in an afternoon. It read token ranges, slices of the ring, from Cassandra, wrote them into ScyllaDB, and recorded each finished range in a local SQLite database so it could resume after a crash. Its estimate was nine days, at up to 3.2 million messages a second.

It stalled at 99.9999% complete. The last few token ranges contained huge runs of tombstones that had never been compacted away, and reading through them timed out, the problem from section 3.6 once again. Compacting that range on Cassandra took seconds, and the copy finished. Discord then sent a small share of reads to both databases and compared the answers, and in May 2022 it made ScyllaDB the primary store.

Cassandra (early 2022)ScyllaDB (after May 2022)
Machines17772
Data per machineabout 4 TB on average9 TB
p99 latency, reading history40 to 125 msabout 15 ms
p99 latency, inserting a message5 to 70 msabout 5 ms

The p99 latency is the time within which 99 of every 100 requests finish, so it describes the slow requests, the ones people notice. Both ranges shrank to a single steady number, and that's what "predictable" meant in Discord's 2017 wish list. In December 2022 the World Cup final between Argentina and France produced a spike in Discord's message traffic for every goal, penalty and half-time, and the graph Discord published of that night is labelled in coalesced messages per second.

Decision

Cassandra was paging engineers at 177 machines. What next?

Keep tuning Cassandra
More machines, more garbage-collector and compaction tuning.
  • No migration risk
  • Hot partitions and GC pauses are built into the design
  • Toil keeps growing with data
A different data model
Move messages to a new kind of store and redesign the queries.
  • Could remove the hot-partition problem entirely
  • Rewrite every query
  • Years of migration
chosen
ScyllaDB plus data services
Same model and keys on a C++ engine; coalesce and route requests in front of it.
  • Tables and queries move unchanged
  • No GC pauses; hot partitions confined to one core
  • Coalescing absorbs bursts before they reach disk
  • Still an LSM store with tombstones and compaction
  • Depends on one vendor's engine

Discord kept the data model it understood and changed the two things around it that hurt most: the runtime underneath (garbage collection) and the traffic shape above (bursts of identical reads). It ended with fewer than half the machines and a p99 that stopped moving. Problems that come from the data model itself, hot current buckets and tombstones from mass deletes, are still there, and the post doesn't claim otherwise.

That settles where Aiko's message lives. It doesn't yet say how it reaches the 320,000 people who are watching #general live, and that path never touches the database at all.

05Fanning a message out to everyone online

5.1A process per guild

Discord's real-time side is written in Elixir, a language that runs on the same virtual machine as Erlang, called the BEAM. Chapter 49's section 4.1 explains why WhatsApp chose Erlang for its connection servers: each connection gets its own process, a tiny independent program inside the virtual machine that costs a few kilobytes, has its own memory, and talks to other processes only by sending them messages. Millions of them run on a handful of CPU cores.

Telephone operators wearing headsets connecting calls with cords at a long switchboard
Operators at a Tokyo telephone exchange. Erlang was created at Ericsson in the 1980s to program the machines that replaced rooms like this one: switches handling thousands of independent calls at once, where one call failing must never take down the rest. Discord's guilds and sessions are the same pattern: many small independent conversations, each with its own process.Photo: NTT Public Corporation, public domain, via Wikimedia Commons

Discord uses two kinds of process. Each connected client has a session process on a gateway machine, which owns that client's WebSocket. Each guild has exactly one guild process, on one of a set of guild machines, which knows the guild's channels, roles and members and which sessions are connected to it. Discord's 2023 post calls the guild process the central routing point for everything that happens in the guild.

So when Aiko's message has been stored, an event saying "new message in #general" is handed to the Starfall guild process. It works out who should see it: every connected member with permission to read #general, and sends a copy to each of their session processes. Each session then pushes it down its WebSocket to the app.

Aiko's message, from send to every screen
ELIXIR ON THE BEAMsendstoreeventfan-outpushAiko's appAPIData services→ ScyllaDBStarfall guild processone per guildSession processesone per client×320kMembers' appsWebSocket each
Step 1. Aiko's app sends the message to the API, which gives it a Snowflake ID.
1 / 5

We've simplified the path from API to guild process here, because Discord hasn't published exactly how the event travels between them. What it has published, in detail, is what happened inside the guild process as guilds grew.

5.2One process can only do one thing at a time

A process handles one message at a time, in order. That's what makes it simple to program, because the guild's state can't change underneath you halfway through, and it's also its limit, because a guild process can never use more than one CPU core.

Discord's 2017 post, written when it reached nearly five million concurrent users, says the original design worked well for guilds of up to about 25 people. By then the /r/Overwatch guild had up to 30,000 people online. Sending one message from one process to another, Erlang's send, cost between 30 and 70 microseconds of the sender's time in Discord's measurements, because the sending process could be paused by the scheduler partway through. At peak, publishing a single event from a big guild took between 900 milliseconds and 2.1 seconds. Meanwhile new events kept arriving, and the guild process's queue of unhandled messages grew. Engineers had to switch off features that generated messages to keep big guilds alive.

Your turn: design it before reading on

Starfall's guild process has to send Aiko's message to 320,000 sessions, at 50 microseconds per send. How long does one message take? And #general is getting 50 messages a second.

5.3Manifold: one send per machine

Discord's first fix, in 2017, was a library it called Manifold, later open-sourced. It starts from a simple observation: Starfall's 320,000 sessions don't live on 320,000 machines. They live on a modest number of gateway machines, each with many thousands of sessions. So instead of sending the event 320,000 times, the guild process groups the destination sessions by the machine they're on, and sends one message to each machine, carrying the event and the list of sessions on that machine. On each receiving machine, a Manifold process splits the list among workers, one per CPU core, and they do the local sends in parallel.

Manifold: from one send per session to one per machine
Guild processStarfallGateway machine 1Manifold partitionerGateway machine 2Manifold partitionerGateway machine 3Manifold partitionermsg→ 320k sessionsmsg+ 110k idsmsg+ 105k idsmsg+ 105k ids
Step 1. Aiko's message reaches the guild process, which has a list of every session that can read #general. Without Manifold it would send to each one in turn, from this one process.
1 / 5

Machine counts and per-machine numbers in the animation are illustrative; a guild this size probably spans many more gateway machines. What Manifold does is real: it takes the per-session cost off the guild process and spreads it over every core of every gateway machine, and it cuts the network traffic between machines too, since each event crosses the network once per machine instead of once per session. In 2017 that was enough for the biggest guilds of the day.

5.4Passive sessions and relays

Manifold makes each send cheaper, but the number of sends still grows with the square of the guild's size. By 2022 the Midjourney guild was heading for a million people online, and Discord needed to remove work, not only spread it.

A first cut came from noticing what most members are doing. Someone in 50 guilds is probably looking at one of them, if any. The other 49 are icons in the sidebar. Their apps need to know very little about those guilds: maybe that there's something unread, and whether they were mentioned. They don't need every message, every typing indicator and every member coming online. Discord made those connections passive: a member who is in a guild but hasn't opened it gets a passive session for it, which receives far less, and becomes active when they click on it. Discord's 2023 post says around 90% of the user–guild connections in large guilds were passive, so the fan-out work dropped by about 90%. Because the work grows with the square of the size, a tenfold cut in work buys only about a threefold increase in how big a guild can be, as the post also points out.

Your turn: design it before reading on

Starfall has 2.4 million members, 320,000 online. If 90% of online connections are passive, how many sessions get every message in #general? And if one machine can handle the full stream for 15,000 sessions, how many are needed?

That 15,000 is a real number. An earlier change, made before Midjourney, had let guilds grow from tens of thousands into the hundreds of thousands. It was to split the guild process's job in two. The guild process keeps the guild's state: channels, roles, members, permissions. Fanning out moves to a set of relay processes that sit between the guild and the sessions. Each relay holds the connections to up to 15,000 sessions and does the fan-out for them, including the permission checks: whether each of its members may see #general. The guild process sends Aiko's message once to each relay, and the relays do the rest in parallel, on different machines.

A big guild with relays and passive sessions
eventfull eventsmall updateAPIStarfall guild processstate: roles, membersRelay 1≤ 15,000 sessionsRelay 2≤ 15,000 sessionsRelay 22≤ 15,000 sessionsActive sessionshave Starfall openPassive sessionsStarfall in sidebar
Step 1. Aiko's message event reaches the Starfall guild process, which now keeps the guild's state but no longer talks to every session.
1 / 4

Relays brought their own surprise. Each relay needs to know about the members it serves, to check their permissions, and at first each relay held a full copy of the guild's member list. For Midjourney that meant dozens of copies of a list of tens of millions of members in memory, and starting a new relay stalled the guild for tens of seconds while it copied the list over. Discord fixed this by sending each relay only the members it needs, a tiny percentage of the whole, according to the post.

Decision

How should one guild reach a million online members?

Guild process sends to every session
The original design: one process, one send per session.
  • Simple: one place knows everything
  • One CPU core does all the work
  • Seconds per message beyond tens of thousands of people
Manifold
Group sessions by machine; one send per machine, split across cores there.
  • Drop-in replacement for send
  • Spreads per-session work over every gateway core
  • The guild still computes everything for every message
chosen
Relays + passive sessions
Relays own up to 15,000 sessions each and do fan-out and permission checks; members not looking at the guild get a cut-down passive stream.
  • Fan-out runs in parallel across many processes
  • About 90% of the work disappears for big guilds
  • More moving parts and more copies of member data
  • The total work still grows with the square of the guild's size

Discord did all three, in that order, each when the previous one ran out. The 2023 post describes the result as scaling individual guilds from tens of thousands to nearly two million concurrent users. It also says these optimisations are the tip of the iceberg: each one moved the bottleneck somewhere new, and finding it was most of the work.

5.5When one mention means checking millions of people

Some events are harder than an ordinary message. When the moderator posted "servers are back, @everyone", the guild had to work out which members can see that channel, all 2.4 million of them, to decide whom to notify. Discord's 2023 post says those checks can take many seconds, and a guild process busy for many seconds isn't handling any other events.

So Discord moved the guild's member data into ETS, a table built into the BEAM that many processes can read at once without copying. Expensive jobs like an @everyone check are handed to separate worker processes that read the member table directly, while the guild process carries on with new events. ETS also helped with something that happens every day: when a guild process moves to another machine, during a deploy or maintenance, its member data, which can be several gigabytes, has to be copied over. With ETS the old guild keeps serving while the copy happens, instead of stalling for minutes.

Even the cleanup of memory caused trouble. Discord tried giving the fan-out step its own process, to take more work off the guild, and performance got much worse. BEAM's garbage collector was to blame: large data passed between processes is tracked separately, and a small threshold made the busy process run full collections, copying a heap gigabytes in size to reclaim a couple of hundred kilobytes. Raising that threshold (an Erlang setting called min_bin_vheap_size) to a few megabytes fixed it, and the change became a clear win. It's the same lesson as Cassandra's pauses in section 3.6, met again in a different runtime: at this scale, how a system frees memory can decide how fast it is.

06The member list

6.1A sorted list that changes every second

On the right of Aiko's window is the member list: members grouped by role (Moderators, then Online, then Offline), sorted by name within each group. In a guild where hundreds of people come online and go offline every second, that list is never still. Now the problem: what does Discord send to Aiko's app so that it can draw it?

Not the whole thing. 2.4 million names, with their roles and statuses, is many megabytes, and it would change hundreds of times a second. But Aiko's window only shows maybe 40 rows at a time. So the app tells the server which part of the list it's looking at, and the server sends just those rows, and after that only the changes that affect them: "insert Ken at position 12", "remove row 30". Discord's 2019 post on this describes it as sending down only the updates for the visible portion of the member list. This is what people call a lazy member list. Exact messages the client sends to ask for a range aren't part of Discord's documented API for bots, so we'll stick to what the post describes.

Aiko's member list: only the visible rows travel
Guild's sorted member list2.4 million entries, in the guildAiko's windowrows 0–39Off screennever sentModeratorsrows 0–5Online A–Brows 6–39Online C–Z~320k rowsOffline~2.1M rowsKen online→ index 12Yuki offlineno update sent
Step 1. The guild keeps the whole list, sorted by role and name. Aiko's app has asked for rows 0 to 39, so only those were sent.
1 / 4

6.2The data structure behind it

zoomDiscordGuild processMember listSorted set with indices

Behind that animation is a demanding data structure. For every join, leave and status change, the guild has to insert or remove an entry in a sorted list of possibly millions, and say at what index it landed, so that clients whose window covers that index can be told. A hash set is no use, because it has no order. A plain sorted array is ordered but inserting at the front means shifting every entry after it. In Elixir, whose data is immutable, it's worse: inserting into a list means building a new list. Discord's 2019 post gives the example of a guild with 100,000 members, where each join built a new 100,001-entry list.

Discord's 2019 post walks through its attempts, measured at 250,000 entries:

StructureWorst insert at 250,000 entries
A sorted Elixir listabout 170,000 µs (170 ms)
Erlang's ordsets (sorted lists)about 27,000 µs
Small sorted lists chained together, skip-list style5,000 µs at the end, 19,000 µs at the front
Chunks that split as they grow640 µs at the end, 4 µs at the front
The same structure in Rust, called from Elixirabout 3.7 µs at worst

What made the difference is the idea in the fourth row. Instead of one long sorted list, keep a sorted list of small sorted chunks. To insert, find the right chunk (a binary search over the chunks' first entries), and insert into that small chunk. When a chunk grows too big, split it in two. The index of an entry is the sizes of the chunks before it plus its position in its own chunk. Every operation touches one small chunk and the short list of chunks, never the whole list.

Finally, Discord wrote it in Rust and call it from Elixir as a NIF (native implemented function), a function compiled to machine code that the BEAM calls directly, built with a library called Rustler. Mutable memory in Rust removed the copying that immutable Elixir data forced, and inserts took between 0.61 and 3.68 microseconds across sizes from 5,000 to a million entries. Discord open-sourced it as SortedSet, and in 2019 it was running every guild's member list, the biggest then having 200,000 members.

?Why not keep the member list in a database?

Because it changes far too often and has to answer instantly. A member coming online changes the list, and in a big guild that happens hundreds of times a second. Each change needs its new index computed and sent to the right clients within moments. Keeping the list in the guild process's memory means a change costs microseconds and no network trip. The durable facts, who is a member and with which roles, are stored in a database; the sorted, live view of them is rebuilt in memory when the guild process starts.

07The gateway connection and compression

7.1One WebSocket, many events

Every event in this chapter so far, Aiko's message, Ken coming online, a member-list update, reaches the app through one connection: the WebSocket between Aiko's app and her session process on a gateway machine. Chapter 49 covers why a persistent connection beats polling, so here are only the details that are specific to Discord, from its developer documentation.

When a client connects, the gateway's first message is Hello, carrying a heartbeat interval (45,000 ms in the documentation's example). From then on the client sends a heartbeat every interval, so both sides notice a dead connection. The first heartbeat waits a random fraction of the interval, so that a million clients reconnecting after an outage don't all heartbeat in the same instant, the reconnection-storm problem from chapter 49's section 4.3. The client then identifies itself, and the gateway sends a Ready event with the state it needs to start: its guilds, channels and settings. A dropped connection can be resumed, replaying missed events without starting again. And every event, for every member, flows down this one connection, so its size matters. On Discord's scale, every byte per event is multiplied by what is probably billions of events a day.

7.2Compressing the whole stream

Discord has compressed this connection since late 2017, at first with zlib, which made messages roughly 2 to 10 times smaller. zlib uses the DEFLATE algorithm, which works in two steps. First it replaces repeated text with short references back to an earlier copy: "the text 40 bytes back, 12 bytes long". Then it encodes the result with a Huffman code. A Huffman code gives short bit patterns to common symbols and longer ones to rare ones.

A binary tree whose leaves are letters with their counts; frequent letters such as space and e sit near the top
A Huffman tree for the sentence 'this is an example of a huffman tree'. Each letter's code is its path from the root, left 0 and right 1, so the frequent space (7 times) gets a short code and the rare x and p get long ones. zlib builds trees like this for the symbols in each block of data.Image: Meteficha, public domain, via Wikimedia Commons

The first step only helps if there's something earlier to point back to. A single MESSAGE_CREATE event is a few hundred bytes of JSON with field names like "channel_id", "author" and "mention_everyone", each appearing once. Compressed on its own, there's little repetition inside it. But the event before it on the same connection had all the same field names, and likely the same channel ID, guild ID and some of the same authors. So Discord compresses the connection as one continuous stream: a single compressor lives for the whole connection, and after each event it flushes, so the event can be sent at once, without forgetting what came before. Each new event can then point back into previous events. zlib keeps up to 32 KB of history to point back into.

In 2024 Discord moved to Zstandard (zstd), a newer compressor from Facebook, released in 2015. Discord first tried it as a dark launch, compressing real traffic with zstd to measure it while users still received zlib. That first test did worse than zlib: an average MESSAGE_CREATE came out over 750 bytes with zstd against about 250 with zlib. Discord's zstd setup was compressing each event on its own, while zlib was streaming. Discord added streaming support to the Elixir zstd library it used and contributed it back, and then MESSAGE_CREATE came out at 166 bytes against zlib's roughly 270, a compression ratio up from about 6 to almost 10, and in under half the time.

You can see both effects with Python 3.14, which has zstd built in. This program makes a thousand events shaped like gateway MESSAGE_CREATE events, in one channel, from a pool of 200 authors, and compresses them four ways: each event on its own and as a stream, with each compressor.

Compressing gateway events one by one versus as a stream
python
Python
import json, random, zlib
from compression import zstd             # Python 3.14+
random.seed(3)
 
words = "gg nice the new model is out who wants to play tonight lol same here".split()
authors = [{"id": str(random.randrange(10**17, 10**18)), "username": f"user{random.randrange(10**5)}",
            "global_name": None, "avatar": "%032x" % random.getrandbits(128)} for _ in range(200)]
 
def message_create(i):                   # a made-up event shaped like a gateway dispatch
    return json.dumps({"t": "MESSAGE_CREATE", "s": i, "op": 0, "d": {
        "id": str(1290000000000000000 + i * 4194304), "channel_id": "1029384756102938475",
        "guild_id": "662267976984297473", "type": 0, "tts": False, "pinned": False,
        "mention_everyone": False, "mentions": [], "attachments": [], "embeds": [],
        "author": random.choice(authors),
        "content": " ".join(random.choices(words, k=random.randint(2, 12)))}}).encode()
 
events = [message_create(i) for i in range(1000)]
raw = sum(map(len, events))
 
def per_message_zlib():                  # a fresh compressor for every event
    return sum(len(zlib.compress(e, 6)) for e in events)
 
def stream_zlib():                       # one compressor for the whole connection
    c = zlib.compressobj(6)
    return sum(len(c.compress(e) + c.flush(zlib.Z_SYNC_FLUSH)) for e in events)
 
def per_message_zstd():
    return sum(len(zstd.compress(e, level=6)) for e in events)
 
def stream_zstd():
    c = zstd.ZstdCompressor(level=6)
    return sum(len(c.compress(e, mode=c.FLUSH_BLOCK)) for e in events)
 
print(f"{'uncompressed':22} {raw / len(events):6.0f} bytes per event")
for name, f in [("zlib, per message", per_message_zlib), ("zstd, per message", per_message_zstd),
                ("zlib-stream", stream_zlib), ("zstd-stream", stream_zstd)]:
    total = f()
    print(f"{name:22} {total / len(events):6.0f} bytes per event   ratio {raw / total:4.1f}")
output
C++
uncompressed              445 bytes per event
zlib, per message         279 bytes per event   ratio  1.6
zstd, per message         300 bytes per event   ratio  1.5
zlib-stream                73 bytes per event   ratio  6.1
zstd-stream                53 bytes per event   ratio  8.3

The two "per message" lines reproduce Discord's surprise from the dark launch: compressed one at a time, neither compressor gets much below 300 bytes, and zstd does slightly worse than zlib. Streaming is the fix. Streaming cuts the size by a factor of four or more for both, because each event mostly points back at earlier ones. And zstd now wins clearly, with a smaller output than zlib's stream. Its much larger window probably explains part of that: an author who last spoke 80 events ago is more than 32 KB back, out of zlib's reach but well inside zstd's. Synthetic events make the exact numbers different from Discord's, but the order of the four lines is the same as in Discord's 2024 post.

Discord tuned zstd's settings, including the window size, and tried training a dictionary, a set of common byte sequences given to both ends in advance. The dictionary shrank the big Ready payload only from 306,745 to 306,098 bytes, made MESSAGE_CREATE slightly worse, and was dropped as not worth the complexity.

7.3Sending passive sessions less

That 2024 project also found a bigger saving in what was being sent, not how it was packed. Passive sessions from section 5.4 received a PASSIVE_UPDATE_V1 event that sent the full current state of the parts of the guild they tracked. Those events were about 2% of all events dispatched but about 35% of all gateway bytes. A new version, PASSIVE_UPDATE_V2, sends only what changed since the last update, and it cut that share to under 5%, about a fifth of all gateway traffic saved. Together, zstd streaming and the smaller passive updates cut the gateway bandwidth used by Discord's clients by almost 40%. Both changes rolled out in the first half of 2024.

08Read states: the unread badge

8.1One record per person per channel

Aiko's message has been stored and delivered. One small piece is left: the white dot next to #general in the sidebar of everyone who hasn't seen it yet, and the red number on any channel where someone mentioned them. Discord calls the record behind this a read state: one per user per channel, with counters such as the number of mentions, updated atomically and often reset to zero. Discord hasn't published the full layout, but probably the simplest way to decide "is anything unread?" is to compare the Snowflake of the last message you read with the newest message in the channel. Snowflakes sort by time, so one comparison answers it.

Discord's 2020 post on the Read States service gives the scale. There are billions of read states. Read states are touched when a client connects, every time a message is sent, and every time someone reads a channel. Each machine keeps an in-memory cache holding tens of millions of read states, least recently used first out, and the fleet handles hundreds of thousands of cache updates a second. Writing each update straight to the database would be far too much, so the cache is the working copy. When a read state changes, a write to Cassandra is scheduled 30 seconds later, and when one is pushed out of the cache, it's written then. That turns hundreds of thousands of updates a second into tens of thousands of database writes, because a read state that changes ten times in 30 seconds is written once.

8.2Garbage collection again: Go to Rust

Discord had written the service in Go, and every two minutes or so its latency and CPU spiked. Go's garbage collector runs at least once every two minutes whether or not memory is short, and with tens of millions of read states in the cache, each collection had to scan all of them, even though almost none had become garbage. Tuning the collector changed nothing visible, and shrinking the cache made the spikes smaller but sent more reads to the database, raising the 99th-percentile latency.

Discord rewrote the service in Rust, where memory is freed the moment its owner is done with it and there's no collector to scan the heap. The first port, finished in May 2019, had no spikes, and after optimising (including switching the cache's map from a hash map to a B-tree map to use less memory) it beat the Go version on every metric Discord tracked. With no collector to scan the heap, a bigger cache no longer meant longer pauses, so Discord gave the machines more memory and raised the cache's capacity again.

Decision

How should a cache of tens of millions of entries manage memory?

A garbage-collected language (Go, Java)
The runtime finds and frees unreachable objects.
  • Easy to write
  • No manual memory bugs
  • The collector scans the whole heap, which is huge
  • Periodic pauses and CPU spikes show up as tail latency
Off-heap memory or a separate cache
Keep the data outside the collected heap, or in Redis or Memcached.
  • The collector sees little
  • Manual serialisation
  • Another hop or another system
chosen
A language without a collector (Rust)
Ownership decides when memory is freed, at compile time.
  • No collector pauses
  • Memory safety checked by the compiler
  • A rewrite
  • A steeper language to learn

Discord's 2020 post is careful to say Go's collector was doing its job; the problem was the workload, a huge heap of long-lived objects that a collector must keep scanning to find almost no garbage. That shape also caused Cassandra's pauses in section 3.6 and the BEAM's in section 5.5. Discord's data services in section 4.3 and its member-list NIF in section 6.2 are Rust for related reasons.

09The whole system

9.1Every box, and why it's there

Discord's message path, end to end
ELIXIR REAL-TIME LAYERsendeventfan-outpushAiko's appMembers' apps320,000 onlineGatewaysessions, zstd-stream×NAPIGuild process+ relays, member listRead statesRust, cache + CassandraData servicesRust, coalescingScyllaDB(channel_id, bucket)
Step 1. Aiko presses Enter. The API gives her message a Snowflake ID, which fixes its place in time and its ten-day bucket.
1 / 7
ComponentWhat it doesAdded because
Snowflake IDs64-bit, time-sorted IDs made without coordinationMessages must sort by time and be created anywhere (§3.2)
(channel_id, bucket) partitionsOne partition per channel per ten daysA channel's history must be together, but bounded (§3.1, §3.3)
ScyllaDBThe message store, a Cassandra-compatible engine in C++GC pauses and toil at 177 Cassandra machines (§4.1, §4.2)
Data servicesRust gRPC layer with routing and request coalescingBursts of identical reads on hot partitions (§4.3)
Guild processOne per guild: state, permissions, routingSomething must know who is in the guild (§5.1)
Manifold, relaysSpread fan-out over machines and processesOne process can't send to 320,000 sessions (§5.2 to §5.4)
Passive sessionsCut-down streams for guilds not being looked atAbout 90% of connections in big guilds are idle (§5.4, §7.3)
Lazy member listOnly visible rows, sorted set with indicesMillions of members, hundreds of changes a second (§6)
Gateway compressionOne zstd stream per connectionHuge numbers of small, repetitive events (§7.2)
Read statesPer-user per-channel counters, cached, written backUpdated on every send and read (§8)

9.2From top to bottom

LevelThe choiceData structure or algorithm
SystemFan-out on read for storage, fan-out on write for live deliveryOne stored copy per message; one pushed copy per active session
IDsTime first in the bitsSnowflake: 42-bit ms timestamp, 5-bit worker, 5-bit process, 12-bit counter
Message storagePartition by channel and time((channel_id, bucket), message_id), bucket = (id >> 22) ÷ 10 days; LSM trees, tombstones, compaction
Storage accessMerge identical readsConsistent-hash routing on channel ID; one in-flight query per key
Fan-outRemove work, then spread itGroup by machine (Manifold); relays of 15,000 sessions; passive sessions
Guild stateShared without copyingETS tables read by worker processes
Member listIndexed sorted setSorted chunks that split; Rust NIF
GatewayCompress the streamDEFLATE (LZ77 + Huffman), then zstd with a large window and per-event flushes
Read statesCache and coalesce writesLRU cache on a B-tree map; delayed write-back after 30 s

10What goes wrong, and what it cost

10.1Failures this design has to survive

What happensWhat the user seesWhat the design does
Everyone opens #general after an @everyoneA short pause, then the channel loadsCoalescing turns thousands of identical reads into a few queries
A bot deletes millions of messages in a channelWithout fixes, a channel that takes 20 s to loadShorter tombstone grace period; empty buckets skipped
A big guild's current bucket gets hammeredSlow requests for that channel, and for others sharing its machinesShard-per-core confines it to one core; coalescing absorbs bursts
A guild grows past what one process can fan outMessages arrive late, then not at allManifold, relays, passive sessions
A guild process moves during a deployWithout care, minutes of stallingMember data in ETS; the old guild keeps serving while it's copied
A gateway machine diesClients reconnectJittered heartbeats and reconnects; resume replays missed events
A garbage collector pausesLatency spikesScyllaDB and Rust services have no collector; BEAM settings tuned

10.2The tradeoffs, in one table

DecisionChosenGiven upWhy it was worth it
Message storageKeep everything, one copy per channelSimple bounded storageHistory is the product: anyone can scroll back years
Partition key(channel_id, ten-day bucket)Even load across machinesThe common read is one sequential scan
DatabaseScyllaDB (2022) over CassandraA widely used Java engineNo GC pauses; 72 machines instead of 177; steady p99
Hot readsCoalesce in a Rust layerOne more service hopThe database stops seeing popularity
Live deliveryPush to every active sessionCheap rooms of unlimited sizeReal-time chat needs it; passive sessions cut 90%
Member listIn memory, only the visible window sentDurability of the sorted viewMicrosecond updates; it can be rebuilt from stored membership
Gatewayzstd stream per connectionMemory per connection for compressor stateSmaller and faster than zlib; with passive deltas, almost 40% less bandwidth
Read statesWrite back after 30 sA few seconds of unread state lost on a crashFar fewer database writes

11Summary

  1. Discord stores one copy of each message, in its channel, because copying into millions of members' mailboxes would multiply every message by the guild's size.
  2. Fan-out to online members grows with the square of a guild's size: more people send more messages, and each goes to more people.
  3. A single indexed database breaks when the data outgrows memory, because Discord's reads are spread almost randomly across all of history.
  4. Snowflake IDs put a millisecond timestamp in the top bits, so they sort by time and can be made anywhere without coordination.
  5. Messages are partitioned by (channel_id, bucket), with ten-day buckets computed from the Snowflake, which keeps reads sequential and partitions bounded, at the cost of hot current buckets and multi-bucket reads for quiet channels.
  6. In an LSM store a delete is a tombstone, and reads must step over tombstones until compaction removes them, which once made a one-message channel take 20 seconds to load.
  7. ScyllaDB kept Cassandra's model without a garbage collector, cutting 177 machines to 72 and p99 reads from 40–125 ms to about 15 ms.
  8. Request coalescing behind consistent-hash routing turns a burst of identical reads into one query per few milliseconds.
  9. A guild process can only use one core, so fan-out moved out of it: Manifold, then relays of 15,000 sessions, plus passive sessions that cut about 90% of the work.
  10. Lazy member lists send only the visible rows, backed by a sorted set of chunks that reports insertion indices in microseconds.
  11. The gateway compresses each connection as one stream, and zstd's larger window beat zlib once it streamed too; smaller passive updates saved even more.
  12. Read states live in a cache with delayed write-back, and moving that service from Go to Rust removed garbage-collection spikes.

12Build this

A tiny Discord message store and fan-out.

  • Run ScyllaDB or Cassandra in Docker and create the ((channel_id, bucket), message_id) table. Generate Snowflakes with the layout from section 3.2 and write a function that fetches the latest 50 messages by walking buckets backwards. Fill one busy channel and one quiet one and count the partitions each read touches.
  • Delete 100,000 messages from one channel and time the read before and after. Run nodetool compact and time it again. Look at the tracing output (TRACING ON in cqlsh) to see the tombstone count.
  • Put a small service in front with request coalescing keyed by channel ID. Fire 10,000 concurrent reads at one channel and compare the database's query count with and without coalescing.
  • Write a fan-out simulator: one "guild" process and N "session" queues. Measure time per message as N grows, then group sessions into relays of 1,000 and mark 90% of sessions passive, and plot the difference.
  • Compress a stream of your own JSON events with zlib.compressobj and compression.zstd.ZstdCompressor, flushing after each, and watch how the ratio changes with the number of distinct authors.

13Interview questions

beginnerWhy does Discord keep one copy of each message instead of one per member, the way WhatsApp groups work?›

A WhatsApp group has at most 1,024 members and deletes messages after delivery, so a copy per member's mailbox is affordable and short-lived. A Discord guild can have millions of members, keeps every message forever, and lets anyone scroll back through history, so a copy per member would multiply storage and writes by the guild's size. Storing one copy in the channel and having readers fetch it (fan-out on read) keeps storage proportional to messages sent. Live delivery to online members is still a push per active session, which is the expensive part Discord spent years scaling.

beginnerWhat's in a Discord Snowflake ID, and why is that useful?›

A 64-bit number whose top 42 bits are milliseconds since the start of 2015, then 5 bits of worker ID, 5 bits of process ID and a 12-bit counter. Sorting IDs sorts by creation time, so a message ID can be the clustering key that orders a channel's messages. Each process can make IDs without asking anyone, since its worker and process bits make them unique. And anyone holding an ID can compute when it was made, which is how Discord derives a message's ten-day bucket from its ID.

intermediateWhy did Discord add a time bucket to the partition key, and what does it cost?›

With channel_id alone, a channel's partition grows forever, and some were over 100 MB in 2016; a partition can't be split across machines and huge ones make compaction and repair slow. Adding a ten-day bucket bounds each partition by ten days of traffic while keeping the common read, the latest messages, a sequential scan of one partition. The costs: a quiet channel may need to read many buckets to find 50 messages (Discord skips known-empty ones), and the current bucket of a very busy channel is still a hot partition, which grew as guilds grew.

intermediateWhat is a tombstone, and how can deleting data make reads slower?›

In an LSM-tree store like Cassandra or ScyllaDB, files are immutable, so a delete writes a marker, a tombstone, that hides the old row until compaction removes both after a grace period long enough for repair to spread the delete to every replica. Reads have to step over tombstones to find live rows. After millions of messages were deleted from one channel, reading its latest messages meant loading millions of tombstones, which took 20 seconds and drove the JVM into long garbage-collection pauses. Discord shortened the grace period from ten days to two, which was safe because repairs ran nightly, and skipped empty buckets.

deepA guild has a million members online. Walk through how one message reaches them, and where the bottlenecks are.›

The message is stored once, in its channel's current bucket, through a data service. An event goes to the guild's single process, which knows the guild's members and permissions. A single process uses one core, and sending to a million sessions at tens of microseconds each takes far longer than the gap between messages, so the work must leave the process. Discord used Manifold to send once per gateway machine and split work across cores there, then relays that each own up to 15,000 sessions and do permission checks and fan-out in parallel, and passive sessions so that the roughly 90% of members not viewing the guild get small deltas instead of every event. Remaining hot spots are whole-guild operations like @everyone permission checks, which run on worker processes reading member data from ETS, and memory management in the BEAM.

deepHow does request coalescing work, and why does it need consistent hashing?›

A coalescing layer keeps a map from request key to the query already in flight. A request for a key with nothing in flight starts a query; requests that arrive while it runs wait on it and get the same result. Identical requests only meet if they reach the same instance, so callers route by a key, the channel ID for messages, using consistent hashing, which also limits how many keys move when instances come and go. The saving grows with how hot a key is: the database sees at most one query per key per query duration, however many readers there are.

14Go deeper

check yourself
A Snowflake's top 42 bits are 41,944,705,796. Roughly when was it made, and which ten-day bucket is it in?›

41,944,705,796 ms after the start of 2015 is about 485 days later, at the end of April 2016. Dividing by 864,000,000 ms (ten days) gives bucket 48, as in the TryIt in section 3.3.

Why did cutting the tombstone grace period from ten days to two not risk deleted messages coming back?›

A tombstone has to outlive the time it takes for the delete to reach every replica. Discord ran repair on the message cluster every night, so every copy had the tombstone well within two days.

Passive sessions removed about 90% of fan-out work in big guilds. Why did that buy only about a threefold increase in maximum guild size?›

Because the work grows with the square of the number of people online. A tenfold reduction in work lets the size grow by the square root of ten, about 3.2 times, before the work is back where it was.

'How Discord Stores Billions of Messages' (Discord blog, 2017)

MongoDB's limits, Cassandra's data model, Snowflake buckets, the upsert race, and the tombstone outage.

'How Discord Stores Trillions of Messages' (Discord blog, 2023)

177 Cassandra machines, ScyllaDB, the Rust data services with request coalescing, the nine-day migration, and the World Cup final.

'How Discord Scaled Elixir to 5,000,000 Concurrent Users' (Discord blog, 2017)

Guild and session processes, the cost of send, Manifold, FastGlobal and the Semaphore library.

'Maxjourney: Pushing Discord's Limits with a Million+ Online Users in a Single Server' (Discord blog, 2023)

Quadratic fan-out, passive sessions, relays, ETS worker processes and the garbage-collection surprise.

'Using Rust to Scale Elixir for 11 Million Concurrent Users' (Discord blog, 2019)

The member list's sorted set, every data structure tried, and the Rust NIF.

'How Discord Reduced Websocket Traffic by 40%' (Discord blog, 2024)

zlib to zstd, why streaming mattered, dictionaries, and passive update deltas.

'Why Discord is Switching from Go to Rust' (Discord blog, 2020)

The Read States service, Go's two-minute garbage collections, and the Rust rewrite.

Discord developer documentation: Reference and Gateway

The Snowflake bit layout, heartbeats, resume, and the zlib-stream and zstd-stream transport compression options.

Designing WhatsApp

Persistent connections, Erlang processes, store-and-forward and fan-out on write for small groups. Chapter 49.

Designing Zoom

Voice and video in rooms, through selective forwarding units, which is how Discord's voice channels work too. Chapter 53.

Storage Engines

LSM trees, SSTables, compaction and why deletes become tombstones. Chapter 18.

Partitioning

Partition keys, hot spots and consistent hashing. Chapter 29.

Replication & Consistency

Quorums, read repair and last-write-wins, the rules behind Discord's upsert race. Chapter 28.

Concurrency Models

Actors and message passing, the model behind Erlang and Elixir. Chapter 15.