Skip to content

Design Discord's Message Storage, stage 6 of 9: change it

Everyone opens the same channel

The cluster has plenty of total capacity. One partition is getting far more reads than three nodes can serve.

System so far· 5 parts
12345CLIENTMembersSERVICEAPI serversSERVICEGatewaySERVICEMessagedata serviceDATABASEMessage cluster

Select a component to see what it is responsible for and which state it owns.

  1. 1Members → API servers: Send, load history, jump
  2. 2API servers → Gateway: New message event
  3. 3Gateway → Members: Push to online members
  4. 4API servers → Message data service: Query by channel (hash-routed)
  5. 5Message data service → Message cluster: Read and write one partition
  • Request / response
  • Asynchronous
  • Server push

What you need to know

0 of 1 checks done
  1. A partition lives on its replicas, typically three nodes. Adding nodes adds capacity for other partitions; it can't split one partition's load. A burst of reads for one channel lands on the same three nodes however big the cluster is.

  2. But hundreds of thousands of "latest page of channel X" requests are the same question. If they meet in one place, they can be answered by one query. That's request coalescing: the first request runs the query; identical requests that arrive while it's running wait for its result. See Request coalescing.

    Requests only meet if they're routed to the same place, which is what consistent hashing on the channel ID does for a fleet of data-service instances.

  3. The same arithmetic applies to a burst of identical reads: compare no protection, coalescing per instance, and one query fleet-wide.

    A hot key expires. Rebuilding it takes one expensive query. Until that query finishes, every request for the key misses. The database can run about 200 of these queries at once before it slows down.
    5,000/s400 ms

    Queries reaching the database while the key is rebuilt (log scale). The dashed line is what it can run at once.

    • No protection2,000
    • Coalesce per server20
    • One rebuild, fleet-wide1
    • Serve stale, refresh behind1
    No protection:
    Every request that misses queries the database. 2,000 requests wait about 400 ms, longer, because the database is past capacity and every query slows down.
    Coalesce per server:
    Each of 20 servers lets one request rebuild; the rest on that server wait for it. 2,000 requests wait about 400 ms.
    One rebuild, fleet-wide:
    A short lock in the cache lets one request in the whole fleet rebuild. 2,000 requests wait about 400 ms.
    Serve stale, refresh behind:
    Keep serving the old value while one background request rebuilds it. Nobody waits; a few requests see a value that is a moment old.
  4. Check

    Why route each channel's requests to one data-service instance?