Design Discord's Trillions of Messages Storage — Cassandra to ScyllaDB Architecture
1. Problem Statement & Scope Clarification
System Mission
Design Discord's distributed real-time messaging storage engine capable of persisting trillions of chat messages across millions of guilds (servers) and channels. The platform must maintain consistent sub-5ms read latencies, handle extreme write bursts ( messages/sec globally during major gaming events or bot activity), prevent hot-partition starvation from mega-channels ( concurrent users in a single channel), and eliminate JVM Garbage Collection stop-the-world pauses through a C++ Seastar shared-nothing thread-per-core architecture.
The Evolution Invariant: MongoDB Cassandra ScyllaDB: Discord originally stored messages in MongoDB, but quickly exhausted RAM and disk capacity as data exceeded single-replica limits. Migrating to Apache Cassandra solved initial horizontal scale, but Cassandra's JVM garbage collection created unpredictable 2-second to 5-second tail-latency spikes (), and tombstones from message deletions severely degraded read performance. Migrating to ScyllaDB (a C++ rewrite of Cassandra utilizing the Seastar asynchronous shared-nothing framework) paired with a Rust Message Data Service for request coalescing eliminated JVM pauses, reduced cluster size by , and locked read latency below .
Synthesizing vector architecture diagram...
Functional Requirements
- Chat Message Persistence (
SendMessage): Ingest messages with strict chronological ordering guaranteed per channel using 64-bit Snowflake IDs. - Channel Message History (
FetchMessages): Retrieve paginated historical windows ( messages before/after a given message ID) with . - Atomic Message Mutations (
EditMessage&DeleteMessage): Update text content and execute soft deletions without triggering tombstone scan storms. - Hot-Partition Isolation: Isolate mega-channels and bot flooding (e.g., Midjourney, gaming announcement channels) so they cannot degrade neighboring channels hosted on the same physical storage nodes.
- Cold Message Archival: Seamlessly transition historical messages older than 2 years to Amazon S3 while preserving queryability.
Non-Functional Requirements (SLAs & SLOs)
- High Availability: read and write availability across multi-AZ AWS deployments.
- Ultra-Low Latency:
- Write Ingestion: , .
- History Read ( messages): , (zero GC-induced pauses ).
- Throughput Capacity: Sustain peak ingestion and peak history lookups.
- Data Durability: Zero data loss ( for acknowledged writes, with
LOCAL_QUORUM). - CAP / PACELC Classification: AP system under CAP; PA/EL under PACELC (Partition Availability; Else Latency over Consistency). Chat history privileges real-time read/write availability over linearizable global order.
Out-of-Scope
- Voice/Video WebRTC peer-to-peer media routing (handled by Discord voice servers).
- Real-time user presence status counters (handled by separate in-memory Redis clusters).
- Content moderation machine learning pipelines (processed asynchronously via Kafka).
2. Capacity & Scale Estimation (Back-of-the-Envelope Math)
Traffic Scale Calculations
- Active Guilds (Servers): active Discord servers.
- Active Text Channels: channels.
- Peak Write Throughput (Ingestion): The 5M/s figure is a short burst ceiling (raids, bot floods, tournament finales); the sustained average is messages/sec: Five years of that growth () on top of the 1 T baseline is what yields the 3 T horizon used below. Peak-to-average ratios of are exactly why the write path needs slowmode at the gateway and the Rust write buffer in front of ScyllaDB.
- Peak Read Throughput (Message History & Polling): Average read-to-write ratio of (channel joins, desktop reloads, scrolling up in history):
Storage Sizing & Cluster Growth (5-Year Horizon)
- Total Persisted Historical Messages: () active baseline; growing to over 5 years.
- Average Message Size:
message_id(Snowflake): 8 byteschannel_id: 8 bytesauthor_id: 8 bytescontent(UTF-8 text): 150 bytes averageembeds& metadata overhead: 60 bytes- Indexing & row metadata overhead: 16 bytes
- Total per message record: .
- Raw Storage Calculation (1 Trillion Messages):
- Replication Factor ( Across Multi-AZ):
- 5-Year Growth Volume (3 Trillion Messages):
Network Bandwidth Calculations
- Write Ingress Bandwidth:
- Read Egress Bandwidth (History fetches averaging 50 messages ): Assuming of reads fetch full 50-message windows ():
Cache & Working Set RAM Calculation
- Active Channel Working Set: Active working set comprises messages sent across the top of active channels within the past 2 hours ().
- Rust Data Service Singleflight Buffers: reads coalesced across 120 Rust proxy instances: aggregate cache footprint.
- ScyllaDB Row & Key Cache: globally.
Fleet Sizing & Compute Infrastructure
- ScyllaDB Nodes (AWS EC2
i3en.12xlarge— 48 vCPU, 384 GB RAM, NVMe SSD ): Usable storage per node at disk high-water mark: . (Throughput governs sizing today: deploy 102i3en.12xlargeinstances. At the 5-year 2.25 PB mark storage catches up: nodes, so the ring grows to ~110 by year 5 regardless of QPS.)
3. AWS-First High-Level Architecture
Synthesizing vector architecture diagram...
Follow a message from send to storage. The client keeps a WebSocket to the Elixir gateway cluster for real-time events and sends messages through the Rust API fleet. In the "Rust Message Data Service Layer" panel, every database call passes through a data service that routes requests for the same channel to the same instance: if 1,000 users open a busy channel at once, the singleflight coalescer turns 1,000 identical reads into one database query, and writes are batched briefly in memory. Only then do requests reach ScyllaDB in the "ScyllaDB Enterprise Storage Tier" panel, a ring spread across three AZs. Message events also stream through Kinesis and Firehose into S3 for archiving. The data service layer is what protects the database from hot channels, the main cause of Discord's earlier Cassandra problems.
Data Flow Walkthrough
- Message Write Ingestion (
SendMessage):- The user types a message and hits enter. The client sends an HTTP POST
/v1/channels/{channel_id}/messagesvia the NLB to theMessage API Service. - The API generates an authoritative Discord Snowflake ID (64-bit integer containing a 41-bit timestamp, 10-bit worker ID, and 12-bit sequence counter).
- The message payload forwards over internal gRPC to the Rust Message Data Service.
- The Rust Data Service computes the bucket index based on the Snowflake timestamp: (10-day time buckets).
- The write is executed asynchronously via the CQL binary protocol to the ScyllaDB cluster using consistency level
LOCAL_QUORUM. - The Rust Data Service stages the newly written message into its local in-memory LRU cache and publishes the event to the Elixir Gateway cluster for real-time WebSocket broadcast to channel subscribers.
- The user types a message and hits enter. The client sends an HTTP POST
- Message History Read (
FetchMessages):- A user opens a channel. The client requests the last 50 messages before a specific Snowflake ID.
- The request hits the
Rust Message Data Service. - Singleflight Coalescing: If 500 users in a gaming guild open the channel simultaneously, the singleflight engine detects identical queries
(channel_id, bucket, before_id)and coalesces them into a single in-flight CQL query. - If the request hits the local in-memory LRU cache, it returns in .
- On a cache miss, the Rust service queries ScyllaDB with
LOCAL_QUORUM. ScyllaDB's core-pinned Seastar thread reads directly from NVMe SSD via Linux AIO/DIO without going through the Linux page cache, returning 50 sorted rows in .
End-to-End Request Tracing Walkthrough
| Step # | Event / Action | Component State | Distributed Transition | Output / Response |
|---|---|---|---|---|
| Step 1 | Client invokes SendMessage API | Client socket active | TLS 1.3 connection terminated at AWS NLB; payload routed to Message API | HTTP POST /v1/channels/9201/messages received |
| Step 2 | Snowflake ID Generation | API Service processing | High-precision 64-bit Snowflake allocated: {timestamp_ms, worker_id, seq} | Message assigned id = 1549890308517425161 |
| Step 3 | Time Bucket Calculation | API Rust Data Service | Rust Data Service extracts epoch timestamp and calculates 10-day bucket | Partition key formed: ((channel_id: 9201, bucket: 427), message_id) |
| Step 4 | In-Memory Write Buffer Stage | Rust Data Service active | Message staged into in-memory ring buffer; immediate LRU cache update | Read-your-own-writes guaranteed locally () |
| Step 5 | CQL Binary Protocol Dispatch | Rust Data Svc ScyllaDB | Request dispatched over core-pinned CQL connection to primary replica node | ScyllaDB coordinator node receives CQL INSERT |
| Step 6 | Seastar CommitLog & MemTable Write | ScyllaDB Core Worker | Thread writes to disk CommitLog via Direct I/O and inserts into in-memory Memtable | Data committed across replicas (LOCAL_QUORUM) |
| Step 7 | CQL Acknowledgment Returned | ScyllaDB Rust Data Svc | CQL ACK received in ; in-flight promise resolved | HTTP 200 OK returned to client via API Service |
| Step 8 | WebSocket Gateway Fanout | Gateway Broker active | Message event pushed to Elixir BEAM pub/sub registry | Event broadcast over WebSockets to 45,000 online guild members |
| Step 9 | Concurrent User Opens History | 250 clients scroll up | Multiple identical FetchMessages(channel: 9201, before: 15498...) queries hit API | Requests arrive at Rust Data Service simultaneously |
| Step 10 | Singleflight Request Coalescing | Rust Data Service engine | Singleflight map binds 250 requests to single pending asynchronous future | Only 1 single CQL SELECT query dispatched to ScyllaDB |
| Step 11 | ScyllaDB Direct I/O Read | ScyllaDB Seastar core | Core reads sorted clustering slice from Memtable + Row Cache + SSTable | 50 message rows returned in |
| Step 12 | Broadcast Result to Callers | Rust Data Service resolving | Future completes; identical result set cloned and distributed to all 250 callers | All 250 HTTP responses return simultaneously () |
4. API Interface Design & Wire Protocol
1. Send Message gRPC Specification (SendMessage)
High-throughput internal RPC between Message API and Rust Data Service.
protobufsyntax = "proto3"; package discord.messages.v1; service MessageStoreService { rpc SendMessage (SendMessageRequest) returns (SendMessageResponse); rpc FetchMessages (FetchMessagesRequest) returns (FetchMessagesResponse); rpc EditMessage (EditMessageRequest) returns (EditMessageResponse); rpc DeleteMessage (DeleteMessageRequest) returns (DeleteMessageResponse); } message MessageRecord { int64 message_id = 1; // 64-bit Discord Snowflake int64 channel_id = 2; int64 author_id = 3; string content = 4; int64 created_at_epoch_ms = 5; int64 edited_at_epoch_ms = 6; bool is_pinned = 7; repeated string attachment_urls = 8; } message SendMessageRequest { int64 channel_id = 1; int64 author_id = 2; string content = 3; repeated string attachment_urls = 4; string idempotency_key = 5; } message SendMessageResponse { MessageRecord message = 1; int32 bucket = 2; } message FetchMessagesRequest { int64 channel_id = 1; int64 before_message_id = 2; // Fetch messages older than this Snowflake int64 after_message_id = 3; // Fetch messages newer than this Snowflake int32 limit = 4; // Default 50, Maximum 100 } message FetchMessagesResponse { repeated MessageRecord messages = 1; bool has_more = 2; }
2. Client REST Wire Contract (GET /v1/channels/{channel_id}/messages)
httpGET /v1/channels/9482019481029384/messages?before=1549890308517425161&limit=50 HTTP/1.1 Host: discord.com Authorization: Bot MTk4MjM0OTAxMjM... X-Discord-Locale: en-US User-Agent: Discord-Android/210.0 HTTP/1.1 200 OK Content-Type: application/json X-RateLimit-Limit: 50 X-RateLimit-Remaining: 49 X-RateLimit-Reset-After: 0.200 [ { "id": "1549890119773745153", "channel_id": "9482019481029384", "author": { "id": "2837192847102938", "username": "GamerOne", "avatar": "a_98fbc10293" }, "content": "GG everyone, outstanding raid!", "timestamp": "2026-09-16T21:10:00.120Z", "edited_timestamp": null, "tts": false, "mention_everyone": false, "pinned": false, "attachments": [] } ]
3. Status Codes & Error Contracts
| HTTP Status | Reason Code | Error Contract Payload | Client Handling / Recovery Strategy |
|---|---|---|---|
200 OK | SUCCESS | Standard message array payload | Render chat timeline immediately |
400 Bad Request | INVALID_SNOWFLAKE | {"error": "Invalid before Snowflake ID"} | Discard cursor; reload from newest channel head |
401 Unauthorized | INVALID_TOKEN | {"error": "Authentication token expired"} | Re-authenticate client session |
403 Forbidden | MISSING_PERMISSIONS | {"error": "User lacks READ_MESSAGE_HISTORY"} | Display lock icon on private channel |
404 Not Found | CHANNEL_NOT_FOUND | {"error": "Unknown channel ID"} | Remove channel from client sidebar |
429 Too Many Requests | RATE_LIMITED | {"error": "Channel rate limit exceeded", "retry_after": 0.5} | Honor retry_after header via client queue backoff |
503 Service Unavailable | SCYLLA_OVERLOAD | {"error": "Storage coordinator timeout"} | Retry with exponential backoff and jitter |
5. Data Models & Storage Architecture
Why Cassandra Failed at Scale and Why ScyllaDB Succeeded
- JVM Garbage Collection (GC) Stop-the-World Pauses: Cassandra is written in Java. When millions of message objects are allocated and discarded during write bursts, the JVM garbage collector (even with ZGC or G1GC) triggers periodic GC pauses ranging from to . This caused catastrophic P99 tail latency spikes and false-positive node gossip timeouts.
- Seastar C++ Shared-Nothing Architecture: ScyllaDB is written in C++ using the Seastar framework. It employs a thread-per-core architecture: each CPU core owns a completely isolated slice of memory, network queues, and NVMe disk channels. Zero cross-core thread locking or lock contention exists. Direct I/O (
O_DIRECT) bypasses the OS page cache entirely, delivering predictable sub-5ms P99 latencies under full saturation. - Partition Bucketing Strategy: Storing all messages for a channel under
PRIMARY KEY (channel_id, message_id)causes unbounded partition growth. Highly active channels accumulate hundreds of millions of messages, creating multi-gigabyte partitions that cause Cassandra compaction stalls, read timeouts, and node crashes. Discord solved this via Time Bucketing.
Synthesizing vector architecture diagram...
Read the relationships top-down: a GUILD (server) contains CHANNELs, and each channel stores MESSAGEs written by USERs. The key design is in MESSAGE: its partition key is (channel_id, bucket), where bucket is a fixed time window, and message_id sorts messages inside the partition. Without the bucket, a busy channel's partition would grow forever, and huge partitions are slow to read, compact and repair; with it, each partition covers one time window and stays bounded. Loading a channel reads the newest bucket first and moves to older buckets only if the user scrolls back. Because message_id is a time-ordered Snowflake ID, messages come back already in time order.
ScyllaDB Production Schema with Time Bucketing
sqlCREATE KEYSPACE discord_chat WITH replication = { 'class': 'NetworkTopologyStrategy', 'us-east-1': 3 }; CREATE TABLE discord_chat.messages ( channel_id bigint, bucket int, -- 10-day time bucket derived from Snowflake ID message_id bigint, -- Snowflake ID: sorts chronologically descending author_id bigint, content text, created_at_epoch_ms bigint, edited_at_epoch_ms bigint, is_deleted boolean, has_attachments boolean, attachment_metadata text, PRIMARY KEY ((channel_id, bucket), message_id) ) WITH CLUSTERING ORDER BY (message_id DESC) AND compaction = { 'class': 'TimeWindowCompactionStrategy', 'compaction_window_unit': 'DAYS', 'compaction_window_size': 10, 'timestamp_resolution': 'MILLISECONDS' } AND default_time_to_live = 0; -- Retain forever in hot tier
Bucketing Formula & Partition Sizing
The bucket index is computed deterministically from the Snowflake ID:
The bucket is computed on the epoch-relative value directly (Discord's make_bucket(snowflake) = (snowflake >> 22) / BUCKET_SIZE); adding the Unix epoch back in would still be deterministic but would not match the value the Rust Data Service derives in §3 and the matrix in §6.4. A message sent on 2026-09-16 lands in bucket 427.
- Partition Size Invariant: In a high-traffic channel generating messages per day, a 10-day bucket caps the partition at rows (). This stays well below the recommended ScyllaDB partition limit, eliminating wide-partition compaction stalls.
Unlock Complete Architecture & Production Runbooks
You have explored the free architectural preview (~44%). Spend 1 Coin to unlock the remaining 6 production deep-dive sections for a full 24 hours.