Skip to content

Request coalescing

When many callers ask for the same thing at the same moment, do the expensive work once and give every caller the result. The fix for thundering herds and hot keys.

Performance & scale

Learn it

0 of 1 checks done
  1. Popular data is requested in bursts: a message in a huge channel, a cache entry that just expired. If every request goes to the database independently, a thousand concurrent readers become a thousand identical queries, and the database falls over computing the same answer a thousand times.

  2. Request coalescing keeps a table of in-flight requests keyed by what's being fetched. The first caller for a key starts the fetch; later callers for the same key wait on that result instead of starting their own. When it completes, every waiter gets the answer and the entry is removed.

  3. 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

    Coalescing runs inside each of 50 instances, and requests are spread randomly. A burst for one key arrives. How many database queries, at most?

Quick reference

The same ideas, condensed for revision.

How it goes wrong

Coalescing on the wrong key
Per-user results are shared between users, leaking data.
A slow leader stalls everyone
All waiters inherit the first fetch's latency or failure; time out the shared fetch.
Scattered routing
Requests for one key land on many instances, so each instance coalesces only a fraction.

Instead, consider

Caching with a long TTL
Repeated reads are spread over time rather than simultaneous.
Precomputation
The hot result is predictable and can be pushed before anyone asks.
Rate limiting the source
Protecting the source matters more than serving everyone.

In practice

singleflight (Go)
A library-level in-process coalescer.
Cache leases
The cache grants one refill per key and asks others to retry.
CDN request collapsing
Edges send one origin request per object while others wait.
Data services behind a hash ring
All requests for a key reach one coalescing instance.

It assumes

  • Concurrent callers can accept the same result (the query is not per-caller).
  • Requests for the same key can be routed to the same place.
  • Waiting a few milliseconds for someone else's fetch is acceptable.

Explain it in your own words

Write at least 60 characters (0 so far). Write it as you would say it in a design review. You will compare it against the points a strong answer makes.

Where you practise it

Further reading

Engineers describing it in systems they run.

  • How Discord Stores Trillions of Messages

    Discord · Bo Ingram · Post, Mar 2023

    The same data model six years on: hot partitions, a service layer that merges identical reads, and a migration of the full history to a new database.

  • Scaling Memcache at Facebook

    Meta (Facebook) · Rajesh Nishtala and others · Paper, Apr 2013

    The reference on running a look-aside cache hard: leases, invalidation from the commit log, failover without hammering the database, and consistency across regions.

  • Caching

    Keeping a copy of data closer to where it is used, trading freshness and complexity for speed and reduced load on the source.

  • Consistent hashing

    Mapping keys to nodes so that adding or removing a node moves only a small share of keys, instead of reshuffling almost all of them.

  • Backpressure and capacity

    When work arrives faster than it can be done, something has to give: the queue grows, the producer slows, or work is shed. Choose which on purpose.