Thousands of Reads, One Query: How Discord Saved Its Database After Migrating Trillions of Messages

| | 7 min read

Discord moved its message database from Cassandra to ScyllaDB and reduced the cluster from 177 nodes to 72. Yet the migration did not make hot partitions disappear. When thousands of people request the same channel data at once, the request path matters as much as the database engine.

Request coalescing can combine overlapping reads into one in-flight query and share the result with waiting callers. The headline describes the technique, not a published Discord benchmark: Discord has not reported a fixed thousands-to-one reduction. Its engineering post describes request coalescing, channel-based routing, and the database migration as related but distinct improvements.

In brief

  • Discord’s move from 177 Cassandra nodes to 72 ScyllaDB nodes improved database capacity and tail latency, but hot partitions could still occur.
  • Request coalescing lets concurrent callers for the same result share one in-flight database read.
  • Routing matching channel requests to the same service instance makes local coalescing more effective.
  • Coalescing is not caching, and it only works safely when the request key captures every input that changes the result.

The problem: database growth and hot partitions

Discord’s message cluster grew from 12 Cassandra nodes in 2017 to 177 nodes storing trillions of messages by the beginning of 2022. The team described unpredictable latency, costly maintenance, compaction backlogs, and JVM garbage-collection pauses. It also encountered hot partitions: a channel and time-bucket pair receiving far more traffic than others.

A large community announcement can bring many users to the same channel at once. Their clients may request overlapping message data. If every request triggers an independent read, the database repeats work for the same partition. Discord reported that overloaded nodes fell behind and that latency affected other queries those nodes served.

Cassandra uses a log-structured merge-tree (LSM-tree) storage engine. Writes are recorded in a commit log and an in-memory memtable, then flushed into immutable sorted files called SSTables. Reads may need to reconcile data from memory and relevant SSTables. Indexes and Bloom filters narrow that work, but storage layout, compaction, query shape, and workload still influence read cost. See the Apache Cassandra storage-engine documentation and the original LSM-tree paper.

Solution 1: improve the database layer

Discord migrated its messages database to ScyllaDB, a Cassandra-compatible database implemented in C++. The team reported that the cluster went from 177 Cassandra nodes to 72 ScyllaDB nodes, with 9 TB of disk per ScyllaDB node compared with an average of 4 TB per Cassandra node.

Discord also reported historical-message read p99 latency falling from 40 to 125 ms on Cassandra to 15 ms on ScyllaDB. Message-insert p99 latency moved from 5 to 70 ms on Cassandra to a steady 5 ms. These are Discord’s production results for its hardware, schema, and workload, not a promise of the same gains elsewhere. Read the full Discord engineering account of the migration.

The database change improved performance and operations, but Discord expected hot partitions could still happen. A more efficient engine does not prevent clients from issuing the same read concurrently.

The problem: duplicate reads still reach the database

Imagine thousands of clients asking for the same channel row within a short period. Even if each query is reasonably fast, sending every duplicate to the database multiplies backend work. Under a burst, that can deepen pressure on the partition and the nodes serving it.

The opportunity is to recognize that a matching operation is already in progress. Later callers can wait for that operation and reuse its result instead of starting another identical database query.

Solution 2: coalesce identical in-flight requests

Discord placed data services between its API and database. These services exposed endpoints for database queries and intentionally held no business logic. Discord says their key feature was request coalescing: one worker performed the database read, while concurrent requests for the same row subscribed to that worker and received its result.

The pattern works in four steps:

  1. The first caller for a key starts the database operation.
  2. Concurrent callers with the same key wait on that in-flight operation.
  3. When it completes, the service returns the result to all waiting callers.
  4. A later request runs the operation again unless a separate cache serves it.

This is often called singleflight or duplicate suppression. It is not caching: the result is shared only while the operation is in flight. Go’s golang.org/x/sync/singleflight package documents the same general pattern, although Discord says its data services were written in Rust.

A small Go example

This runnable example simulates loading channel data concurrently. The sleep stands in for a database read.

package main

import (
	"fmt"
	"sync"
	"time"

	"golang.org/x/sync/singleflight"
)

var reads singleflight.Group

func load(channelID string) (any, error) {
	value, err, _ := reads.Do(channelID, func() (any, error) {
		fmt.Println("one database query for", channelID)
		time.Sleep(100 * time.Millisecond)
		return "recent messages", nil
	})
	return value, err
}

func main() {
	var callers sync.WaitGroup

	for i := 0; i < 5; i++ {
		callers.Add(1)
		go func() {
			defer callers.Done()
			result, err := load("channel-42")
			fmt.Println(result, err)
		}()
	}

	callers.Wait()
}

Save it as main.go in a Go module and run go mod init example.com/coalescing, go get golang.org/x/sync/singleflight, and go run .. Output order varies, but overlapping calls for the same key share one function execution. Go 1.18 or newer supports the any type used here.

The problem: matching reads may reach different service instances

An in-memory coalescer coordinates only callers that reach its process. If matching requests are scattered across several service instances, each instance can start its own database call. Coalescing alone therefore does not guarantee one query across a fleet.

Solution 3: route matching database requests together

Discord used consistent hash-based routing so requests with the same channel ID went to the same data-service instance. That made concurrent requests more likely to meet in the same in-flight request group and share one database read.

Routing is not free capacity. A very hot channel can concentrate incoming traffic on one service instance, and membership changes or retries can affect placement. The service still receives and answers each request even when it avoids repeating the database work.

Choose the database coalescing key carefully

The key must represent the full meaning of the operation. If callers request different pages, cursors, limits, or visibility scopes, a channel ID alone may be too broad. Include every input that changes the result, for example:

key := fmt.Sprintf(
    "messages:%s:before=%s:limit=%d:scope=%s",
    channelID, beforeID, limit, visibilityScope,
)

Authorize every caller correctly. A caller must not receive data merely because another caller was allowed to request it. Errors are shared too: if the in-flight operation fails, waiting callers may all observe that failure. Retries, timeouts, and fallback behaviour need their own policies.

What request coalescing solves for the database

Coalescing helps when many callers request equivalent data concurrently. It can reduce repeated database reads during a burst. It does not fix unique queries, inefficient query plans, hot-key skew, compaction backlogs, or storage-engine limitations. It also does not make the leader’s query faster: every waiting caller still waits for it.

A cache may serve later requests over a time window, but then freshness and invalidation become part of the design. Coalescing avoids that decision because it shares only work already in progress. The techniques can complement each other, but they solve different problems.

For a wider foundation in partitioning, replication, caches, and service boundaries, see our system design fundamentals guide.

Takeaway: improve the database and the request path

Discord’s story involves two distinct improvements. ScyllaDB addressed database performance and operational characteristics. Data services reduced duplicate reads before they reached storage, and routing by channel ID helped equivalent requests meet at the same instance.

Before adding coalescing, work through these checks:

  • Do equivalent requests overlap in time?
  • Does the key include every input that changes the result, including visibility and pagination?
  • Will matching traffic reach the same process, and can that instance handle a hot key?
  • What do query counts, wait time, failures, hot-key skew, and tail latency show after rollout?

These measures show whether coalescing reduces database pressure or simply moves the queue into the application tier.

Subscribe to Our Newsletter

We don’t spam! Read our privacy policy for more info.