Discord's engineering problem is a pub/sub system with very uneven topics. Most servers (guilds internally) have a handful of members; a few have millions. Every message, presence change and reaction must reach every connected member who can see it within a fraction of a second, while voice needs latency in tens of milliseconds. So Discord runs three planes: an HTTPS API that accepts writes, a gateway that pushes events, and a voice plane that moves media over UDP.
This article walks through each plane using Discord's developer documentation and engineering blog, and says so where a detail is unpublished. For the generic problem, read designing a real-time chat system first.
The three planes and why they are separate
Writes and reads have different shapes. A message is written once and read thousands of times, so the write path is an ordinary request/response API: the client sends POST /channels/{channel.id}/messages over HTTPS and gets the persisted message back. The read path is push: every client holds one long-lived WebSocket to the gateway. Voice is continuous, loss-tolerant media where a late packet is worthless, so it uses UDP and its own servers.
Each plane scales and fails on its own. A storage slowdown delays sends but does not drop gateway connections, and a gateway restart does not interrupt a voice call.
The gateway protocol, step by step
The gateway is documented for bot developers, and its lifecycle shows the reliability contract Discord offers any client. Every payload carries an opcode op, data d, and for dispatched events a sequence number s and event name t. The lifecycle is:
- Connect to the gateway URL and receive Hello (op 10), which carries
heartbeat_intervalin milliseconds. - Wait
heartbeat_interval * jitter(random 0 to 1) before the first Heartbeat (op 1), then send one every interval with the last sequence number seen; the server answers with Heartbeat ACK (op 11). Jitter stops mass reconnects from heartbeating in lockstep. - Send Identify (op 2) with the token, the intents you want, and optionally a shard pair.
- Receive the Ready dispatch (op 0), which carries a
session_idand aresume_gateway_url. - From then on, receive Dispatch (op 0) events in sequence order.
- On disconnect, reconnect to
resume_gateway_urland send Resume (op 6) with the session id and last sequence; the server replays what you missed. If the server sends Reconnect (op 7) you do the same. If it sends Invalid Session (op 9) you must identify again from scratch.
A missing Heartbeat ACK is the only reliable sign of a zombie connection, so close and resume if no ACK arrives before the next heartbeat is due.
# Python-flavoured pseudocode: connect/spawn/backoff stand in for your WebSocket and task library
PROPS = {"os": "linux", "browser": "mybot", "device": "mybot"}
async def run_gateway(url, token, intents):
session_id, resume_url, seq = None, None, None
while True:
ws = await connect(resume_url or url, compress="zstd-stream")
hello = await ws.recv_json() # op 10
interval = hello["d"]["heartbeat_interval"] / 1000
hb = spawn(heartbeat_loop(ws, interval, lambda: seq)) # first beat after interval * random()
if session_id:
await ws.send_json({"op": 6, "d": {"token": token, "session_id": session_id, "seq": seq}})
else:
await ws.send_json({"op": 2, "d": {"token": token, "intents": intents, "properties": PROPS}})
try:
async for msg in ws:
if msg["op"] == 0: # dispatch
seq = msg["s"]
if msg["t"] == "READY":
session_id = msg["d"]["session_id"]
resume_url = msg["d"]["resume_gateway_url"]
handle(msg["t"], msg["d"])
elif msg["op"] == 7: # server asks us to reconnect and resume
break
elif msg["op"] == 9: # invalid session: start over
session_id, resume_url, seq = None, None, None
break
except (ConnectionClosed, HeartbeatAckTimeout):
pass # fall through and resume
finally:
hb.cancel()
await backoff_with_jitter()The documented limits shape client behaviour. A connection may send at most 120 gateway events per 60 seconds, and an outbound payload above 4,096 bytes closes the connection with code 4002. Close code 4004 means authentication failed, 4009 that the session timed out, 4010 an invalid shard, and 4013 and 4014 invalid or disallowed intents. Never retry 4004 or 4014 in a loop; they fail the same way every time. For bandwidth, JSON can use transport compression (zlib-stream or zstd-stream), where one compression context spans the connection, or ETF, Erlang's binary term format.
Sharding the gateway
A bot in many guilds cannot take every event on one socket, so the gateway shards by guild. The formula is public: shard_id = (guild_id >> 22) % num_shards. The right shift discards the low 22 bits of the guild's Snowflake ID (worker, process and counter) and keeps the creation timestamp, so guilds spread evenly across shards over time. Discord requires sharding once a bot is in 2,500 guilds, and a shard holds at most 2,500. Direct messages always go to shard 0.
Identify calls are rate limited per bucket: max_concurrency (returned by GET /gateway/bot) says how many shards may identify in each 5-second window, and a shard's bucket is shard_id % max_concurrency, so a large bot boots in waves. Keep shard count stable across deploys: changing it remaps every guild.
Sessions and guilds as Erlang processes
Behind the gateway, Discord's real-time core is written in Elixir on the Erlang VM (BEAM). Discord's 2017 post on reaching five million concurrent users describes the model. Each connected client gets a session process, a GenServer that tracks what that user can see. Each guild gets a guild process on some node in the cluster, located by a consistent hash ring. When something happens in a guild, the guild process sends the event to every session process subscribed to it, and each session filters by permission and pushes to its socket.
BEAM processes are cheap, isolated and preemptively scheduled, so one busy guild cannot starve a node. The cost is that a guild process is one sequential mailbox: with a million online members, one event means a million sends from one process.
Discord's published fixes all shrink that work. manifold, its open-source library, groups destinations by node and sends one batch per node. Later write-ups on very large servers add relays, processes that each fan out to a slice of members, and passive sessions, connected members not looking at the guild, who get a thinner event stream. The exact sizing is internal; the transferable pattern is hierarchical fan-out and sending less to people who are not watching.
Snowflake IDs
Every Discord object (user, guild, channel, message) has a 64-bit Snowflake ID. Bits 63 to 22 hold milliseconds since the Discord epoch, 1420070400000 (the first second of 2015). Bits 21 to 17 hold an internal worker ID, bits 16 to 12 a process ID, and bits 11 to 0 a per-process counter. IDs are unique without coordination and sort by creation time, so "messages before X" is a range scan.
DISCORD_EPOCH_MS = 1420070400000
def decode(snowflake: int) -> dict:
return {
"unix_ms": (snowflake >> 22) + DISCORD_EPOCH_MS,
"worker": (snowflake >> 17) & 0x1F,
"process": (snowflake >> 12) & 0x1F,
"counter": snowflake & 0xFFF,
}
def lower_bound_for(unix_ms: int) -> int:
# smallest ID that could have been created at unix_ms: useful for "since" queries
return (unix_ms - DISCORD_EPOCH_MS) << 22Pagination parameters such as before and after take a message ID, so a synthetic ID works as a time cursor. IDs travel as JSON strings because JavaScript numbers lose precision above 2 to the power 53.
Message storage: buckets, Cassandra and ScyllaDB
Messages moved from MongoDB to Cassandra around 2016. The 2017 post on storing billions of messages gives the schema: partition by (channel_id, bucket) and cluster by message_id, where the bucket is a fixed time window derived from the Snowflake. Discord sized it at roughly ten days after finding that ten days of the busiest channels stayed under about 100 MB per partition. Reading a channel's latest messages is a single-partition scan in descending ID order; scrolling back walks to earlier buckets.
CREATE TABLE messages (
channel_id bigint,
bucket int, -- (message_id >> 22) / BUCKET_MS, BUCKET_MS = 10 days
message_id bigint,
author_id bigint,
content text,
PRIMARY KEY ((channel_id, bucket), message_id)
) WITH CLUSTERING ORDER BY (message_id DESC);The same post records two lessons for any wide-column store: writing a null still writes a tombstone, so write only non-null columns; and last-write-wins can turn a concurrent edit and delete into a half-row, so detect rows missing required fields on read and delete them.
By early 2022 the cluster had grown to 177 Cassandra nodes holding trillions of messages, and the 2023 post describes the problems: hot partitions when one channel got very busy, JVM garbage-collection pauses, and compaction backlogs that hurt read latency for everything sharing a node. Discord moved to ScyllaDB, a Cassandra-compatible database in C++ with a shard-per-core design and no garbage collector, and reported 72 nodes and p99 reads of about 15 ms, down from 40 to 125 ms. A comparison of the two engines is in Cassandra vs ScyllaDB.
The database swap was half the fix. The other half was a set of data services in Rust between the API and the database. Requests are routed by channel ID on a consistent hash, so all reads for one channel land on the same instance, and that instance coalesces them: if a thousand clients ask for the same recent messages at once, one query runs and a thousand callers get the result. A hot channel stops being a hot partition. Bucketed time-series modelling in general is covered in Cassandra time-series modelling.
Voice: a selective forwarding unit and end-to-end encryption
A client joining a voice channel tells the main gateway, receives a voice server endpoint and token, opens a separate voice WebSocket for signalling, and sends Opus audio as RTP over UDP. Discord's 2018 post on 2.5 million concurrent voice users describes WebRTC in browsers, a customised WebRTC stack in native clients, and Discord's own C++ media servers.
Those servers are a selective forwarding unit (SFU): they forward each speaker's packets instead of decoding and mixing, which keeps server CPU low and adds no transcoding delay, at the cost of more downstream bandwidth in large calls. The same trade-off shows up in Zoom's architecture.
Since September 2024 Discord has rolled out DAVE, its audio and video end-to-end encryption protocol, for DMs, group DMs, voice channels and Go Live streams (Stage channels were excluded). DAVE uses Messaging Layer Security (MLS) for group key exchange and WebRTC encoded transforms to encrypt each media frame with a per-sender key. Transport encryption to the SFU remains, so the SFU still routes packets but cannot read their media. Discord said DAVE would become required; check its current documentation for enforcement status.
Worked example: one message in a very large guild
Follow a single message posted in a busy channel of a guild with 800,000 members, 120,000 of them online.
- The client posts to the API with a nonce. The API checks permissions and rate limits, assigns a Snowflake, writes through a data service to ScyllaDB and returns the stored message.
- The API publishes
MESSAGE_CREATEto the guild process, which hands slices of the 120,000 sessions to relays rather than messaging each one itself. - Session processes drop the event if their user cannot see the channel; passive members get a thinner stream.
- Each session assigns its next sequence number, compresses and writes to the socket. A client that disconnected 20 seconds ago gets the event during Resume replay.
- When 5,000 people then scroll the channel, the same data service instance coalesces their identical reads into one partition query.
No single process does work proportional to the guild's size.
Failure modes and what absorbs them
| Failure | Symptom | What absorbs it |
|---|---|---|
| Gateway node restart | Thousands of sockets close together | Resume with session id and sequence; jittered heartbeats and backoff |
| Zombie TCP connection | No events, socket looks open | Missing Heartbeat ACK triggers close and resume |
| Huge guild event storm | Guild mailbox grows, events lag | Relays, node-batched sends, passive sessions |
| Hot channel reads | One partition saturates a node | Data-service coalescing and consistent-hash routing |
| GC pause or compaction backlog | p99 read spikes for unrelated channels | GC-free ScyllaDB, shard-per-core isolation |
| Voice packet loss | Audio artefacts | Late packets are dropped, never retried |
| Bot reconnect loop | Repeated 4004 or 4014 closes | Treat as fatal: fix the token or intents instead of retrying |
Trade-offs to copy, and ones to avoid
Copy: a write API separate from a resumable, sequenced push channel; time-sortable IDs; time-bucketed partitions; read coalescing before blaming the database; an SFU for group voice.
Think twice: one process per guild keeps ordering simple but caps per-guild throughput, and lifting that cap took Discord years. If you expect very large rooms, design hierarchical fan-out from the start. A BEAM cluster is also a staffing decision. A broker-based fan-out like the one in publish/subscribe system design trades some latency for more familiar operations. For a contrasting design built around one-to-one end-to-end encryption, see WhatsApp's architecture.
What to do next
- Write a minimal gateway client for a test bot handling Hello, heartbeats, Identify, Resume, op 7 and op 9; cut the network mid-session and confirm Resume replays events.
- Decode a few real Snowflakes from your own messages with the snippet above and check the timestamps.
- Model a bucketed message table for your own chat product: choose a bucket size from your busiest channel's write rate so partitions stay under about 100 MB.
- Prototype request coalescing in front of your hottest read path and measure the database query rate before and after a burst.
- Load-test fan-out: simulate one room with 100,000 subscribers and measure p99 delivery delay with flat versus two-level (relay) fan-out.
- Run an open-source SFU and measure downstream bandwidth per participant at 5, 25 and 100 participants.