Programmable RevenueBuilder Edition Subscribe
Issue 01 Builder Spotlight Reported

The hot-partition problem, and how Discord beat it at trillion-message scale

When one entity gets popular, one database partition takes the load meant for a fleet, and caching alone cannot save it. The clearest public answer is Bo Ingram's work at Discord: shape the traffic in a service you control, then change the engine.

Thomas Cornelius August 22, 2026 7 min read
Technical schematic: dozens of request arrows converging into a coalesce service that sends one query to a database.
N identical requests in flight become one database read.
Engraved line-art portrait of Bo Ingram Photo: graph8
Builder Spotlight Bo Ingram Staff Software Engineer , Discord

Every partitioned datastore has the same soft spot. Data is spread across nodes by key, which works until one key gets popular: one viral post, one busy channel, one enormous account. Then a single partition takes traffic meant for a fleet, its latency climbs, and, because callers retry, the damage cascades outward to requests that never touched the hot key. The failure has a name, hot partition, and it sits under most of the big-system outages that begin with “one entity got popular.”

Builder Spotlight takes one problem like this per profile and points at the builder whose public work is the standout answer. For hot partitions, that is Bo Ingram, staff software engineer at Discord, whose account of moving trillions of messages off a dying Cassandra cluster is the reference story for this class of failure.

The problem, across the board

The textbook responses each stall in a known place. Caching is the first reflex, and it has a hole with its own name: when a hot key expires, every concurrent miss goes to the database at once, the cache stampede. JVM tuning shrinks garbage-collection pauses that then return as the heap grows. Adding nodes spreads the average load and does nothing for the one hot key, since a partition lives on one replica set no matter how many nodes surround it. Swapping to a managed store trades known failure modes for new ones and leaves the stampede unsolved.

At Discord the stalls were not textbook, they were operational. The message cluster had grown from 12 Cassandra nodes to 177 in five years. Hot partitions in busy channels sent latency cascading across the cluster, and garbage-collection pauses got bad enough that nodes needed manual reboots. Ingram described a cluster with, in his words, “serious performance issues that required increasing amounts of effort to just maintain.” The system worked. It consumed the people running it.

What Discord built instead

Two moves, and the order matters.

First, the team put data services written in Rust between the API and the database. The services coalesce requests: when many users ask for the same message at once, the database gets asked once. Routing by a hash of the channel ID means the same channel always lands on the same service instance, so the coalescing actually catches the duplicates. The database stopped seeing traffic spikes raw.

get msg 812, chan 4 get msg 812, chan 4 get msg 812, chan 4 rust data service coalesce identical reads · route by hash(channel id) scylladb: one query hot partition sees 1 read, not N the same layer flattened spikes during the 2022 World Cup final: nine goal-shaped surges, absorbed before the database felt them
Coalescing: N identical requests in flight become one database read.

Second, with traffic already flowing through an interface the team controlled, they swapped the engine underneath it: ScyllaDB, Cassandra-compatible but written in C++ on a shard-per-core design, which removes the garbage-collection class of failure outright. Their first migration plan, using Spark, was estimated at three months. They wrote a purpose-built migrator in Rust and finished in nine days, moving up to 3.2 million messages per second.

The published results: 177 nodes became 72, each carrying 9 TB instead of roughly 4. Historical message reads went from a p99 of 40 to 125 milliseconds to 15. Inserts went from a p99 swinging between 5 and 70 milliseconds to a steady 5. During the 2022 World Cup final, nine traffic spikes tracked the goals on the pitch, and the system absorbed all of them.

Why this was the right technical direction

The insight worth stealing is that the traffic layer and the storage engine are separate problems, and Discord solved them in the correct order. The coalescing tier fixed what no engine swap could: a hot partition is hot because N requests arrive, and only something upstream of the database can turn N into 1. The engine swap fixed what no service tier could: garbage-collection pauses live inside the JVM, and the only way out is an engine without one. Doing traffic first meant the riskiest change of the project, the migration, ran behind a stable interface, invisible to callers.

The migrator is the third, smaller lesson: a general-purpose tool priced the one-shot job at three months, a purpose-built one did it in nine days. Migrations at this scale are exactly the case where building the tool wins.

What applies today

Request coalescing has since become standard practice, with library support in most stacks; Go ships it as singleflight, and the pattern shows up in most serious caching guides. The widely missed detail is that these libraries are per-process: a singleflight group only deduplicates calls inside one application instance. Discord’s version works fleet-wide precisely because the coalescing lives in a dedicated service tier with hash routing, so all requests for a channel converge on one process first. If you deploy singleflight across fifty stateless API pods without that routing, a hot key still hits your database fifty ways. That distinction is the part of this story most current write-ups still get wrong, and it is why the architecture, not the library, is the accomplishment.

Run the same play

The move that made the migration safe was not the new database. It was the layer in front of it, and it is buildable in a week at most scales.

  1. Put a thin service between your hottest reads and the database, and coalesce identical in-flight requests per key (Go’s singleflight, or a futures map keyed by request) so N concurrent identical reads become one query.
  2. Route requests by a hash of a stable key, channel, account, contact, so the same entity always lands on the same service instance and its cache.
  3. Only after traffic is shaped, evaluate the engine. Migrate with a purpose-built copier that walks token ranges and verifies counts, and measure it against the general-purpose tool before committing to either.

That pattern transfers directly to anyone putting a database behind high-volume sending, enrichment, or event traffic, which is the shape of every workload on the graph8 API.

About this section: Builder Spotlight profiles an engineer outside graph8, from their own published work only.

Source: Bo Ingram, “How Discord Stores Trillions of Messages”, Discord engineering blog, March 2023. All Discord numbers are from that post. More from Ingram: his ScyllaDB Summit talk on the migration and ScyllaDB in Action. Background: ScyllaDB’s architecture documentation on shard-per-core, and Go’s singleflight package for the in-process version of coalescing.

Build with graph8