advanced 3 min answer Multiple choice

Discord's message store reached trillions of rows on a Cassandra cluster reported at 177 nodes; in 2023 it published a move to ScyllaDB and placed Rust data services in front of the store whose main job is request coalescing. Messages are partitioned by channel id plus a time bucket. Which property of the workload best explains why the answer was a different wide-column engine plus a coalescing layer rather than a relational sharding tier?

discordscylladbcassandrahot-partitionrequest-coalescing
Pick one
Show the full answer Hide the answer

The deciding property

The query surface was already closed. The message store answers one question: give me a page of messages for this channel in this time range. A relational query planner earns its cost by answering questions you did not design the physical layout for, and here there are none to answer. So the thing a relational tier is good at was worth nothing, and the two things that actually hurt were untouched by it.

Those two things are tail latency from runtime behaviour and duplicate reads of one partition. Cassandra's JVM garbage collection and compaction put the pain in p99 rather than in throughput - the same class of problem Discord described in 2020 when it moved its Read States service from Go to Rust and attributed the periodic latency spikes to the runtime's collector rather than to load. ScyllaDB's C++ shard-per-core design removes a collector from the read path, and because it speaks the Cassandra protocol, the data model and the client code survived the move.

Why coalescing is the part a storage engine cannot do

A popular channel is one partition. When a message lands, thousands of connected clients ask for the same page within milliseconds. That is N identical reads of one partition, and no storage engine can remove them, because the duplication happens upstream of the store. A data service that holds in-flight queries and attaches new identical requests to one outstanding query turns N reads into one and fans the result back out. That is why the Rust services sit between the API and the database rather than beside it.

Why the other options fail

  • "Past what a relational engine can store." Volume is the reason usually given and almost never the reason; sharded relational systems hold petabytes in production. Volume forces partitioning, it does not pick an engine.
  • "Cannot route on a composite key." A sharding proxy routes on whatever key you hand it, and composite routing keys are ordinary in every production sharding layer. The option sounds technical and is simply untrue.
  • "Stronger transactional guarantees." Backwards. ScyllaDB is a Cassandra-compatible engine with comparable guarantees; the migration was chosen partly because nothing about the consistency model had to change. Wanting transactions would have argued for the relational option.
  • "Standardised on Rust." Language choice followed the latency requirement rather than leading it. Treating it as the cause inverts the reasoning, and it is how a design review gets nodded through for the wrong reason.

What would flip the decision

If this changes Choose Because
The product invents new query shapes every month Relational, or a search index beside the store A planner is cheaper than a new physical layout each quarter
Reads are unique per user rather than identical Partitioning and caching, not coalescing Coalescing earns nothing when no two requests match
Tail latency traces to disk rather than runtime Stay on the engine and fix storage A migration will reproduce the problem on new hardware
Volume is tens of billions of rows, not trillions One relational primary with partitioned tables The operational cost of a distributed store is not yet justified

When this is the wrong thing to copy

This outcome rests on three conditions holding together: the data model was already correct, the destination was protocol-compatible so the migration was mechanical, and the team could build and operate a coalescing layer. A team of eight with 50 GB of messages copying only the destination takes on the operational burden and gets none of the benefit. If your symptom is identical concurrent reads of one hot key, build the coalescing layer and keep your current database - that is the part of this story that generalises, and it is the cheaper half.