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
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.
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.
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 msQueries 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.
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?Coalescing collapses simultaneous requests; a cache serves repeated ones. Related tools: stale-while-revalidate serves the old value while one caller refreshes it, and jittered TTLs stop many keys expiring at once. At the cache layer, the same idea appears as a lease on miss: one caller may refill; the rest wait briefly and retry.
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
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.
Related concepts
- 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.