WeChat's backend sheds overload with DAGOR, published at SoCC 2018: each service detects overload from the average queuing time of requests in its own pending queue - thresholded at 20 ms against a 500 ms task timeout - and rejects by admission level rather than by rate. A platform team running a service mesh asks why they cannot get this from the mesh's load-shedding settings. What did WeChat choose, why did it fit, and what does it cost?
Show the full answer Hide the answer
The situation they were in
A very large microservice fleet — the paper describes WeChat's backend as thousands of services — where one user action fans out across many hops. Two properties make rate-based shedding useless there. A service's response time includes its downstream's latency, so shedding on response time makes a healthy service reject traffic because of someone else's overload. And a request rejected after five hops has already consumed the work of five services, which is the load you were trying to remove.
What they chose
Overload is detected locally from queuing time, and shedding is by priority, not by rate.
- Signal: the average time requests spend waiting in the service's own pending queue, with 20 ms as the overload threshold against a 500 ms default task timeout. Queuing time is local, so it measures this service's saturation and nothing else.
- Admission level: each business priority is subdivided into 128 user-priority levels, giving tens of thousands of distinct levels. Business priority comes from the entry-point action; user priority is derived by hashing the user id, so a given user is not always the one dropped.
- Control loop: each server keeps a histogram of the priorities it is receiving and, under overload, moves its admission threshold to the bucket that cuts expected load by about 5%; when not overloaded it relaxes by about 1%. Decentralised, per machine, with no global coordinator.
- Propagation: a server piggybacks its current admission level on the responses it sends upstream, so the caller stops issuing requests below that level before making the call. Work is refused at the cheapest point in the call graph rather than the most expensive.
The paper reports the scheme running in the WeChat backend for more than five years.
Why it fit their constraints
Priority answers the question a rate limit cannot: when you must drop 30% of traffic, which 30%? Sending a message and refreshing a recommendation list are not interchangeable. Propagation is what makes the saving real, because a deep call graph wastes the most work on requests that will be rejected late.
What it cost them
An organisation-wide priority scheme, which is the part a mesh cannot supply. Someone must rank every entry-point action, keep the ranking current as products ship, and have the authority to refuse a team that wants its traffic classified critical. Priority inflation is the failure mode: if most traffic is top priority the admission level has nothing to cut, and the symptom is a shedding system that appears configured and does nothing. The defence is an audit of the priority distribution as a standing metric.
It also costs a framework change. Priority has to travel in the call context and be honoured by every hop, so this is a change to the RPC library and to every service's handler path — not a mesh configuration. A sidecar can enforce a concurrency limit, but it does not know which request matters.
When not to copy it
With a handful of services, one product and a shallow call graph, a concurrency limiter plus a queue timeout at the edge gets most of the benefit for a day's work. DAGOR's value scales with call-graph depth and with the number of distinct products sharing one fleet. If your traffic is one product with uniform business value, there is no priority to express and the taxonomy is pure overhead — and a two-tier split, interactive versus batch, captures nearly all of the available benefit.