advanced 3 min answer

Tencent's DAGOR (SoCC 2018) detects overload in WeChat's microservice fleet from the average queuing time in a server's pending queue rather than from CPU, treats 20 ms as the overload threshold, and piggybacks the server's current admission level onto its responses so callers shed doomed work themselves. A streaming ingest tier has the same shape. What transfers, and where would copying it be a mistake?

tencentwechatload-sheddingadmission-controlingest
Show the full answer Hide the answer

The situation they were in

A user action in WeChat fans out into calls across a large microservice estate. Shedding a request at the fifth hop wastes the work of the first four, and shedding randomly means a user whose session is half written sees a broken product rather than a slower one. CPU utilisation was a poor saturation signal: a server pinned at high CPU that still drains its queue promptly is not overloaded.

What they chose

  • Queuing time as the signal. The average wait in the pending queue, with 20 ms as the threshold, measured over a window of one second or 2,000 requests, whichever comes first. This measures the thing users feel instead of a resource that correlates with it.
  • Two-dimensional priority. A business priority assigned at the entry service and inherited by every downstream request, with 128 user-priority levels underneath it derived from the user identity, so shedding is consistent for a given user rather than random per hop.
  • Collaborative admission control. A server attaches its current admission level to each response, and the caller applies it locally before sending the next request, so work destined to be rejected is dropped before it crosses the network.

What transfers to a streaming ingest tier

  • Queuing delay at the ingest gateway is the right saturation signal. Accept-queue wait grows while CPU stays flat, because the bottleneck is usually the downstream produce call or a broker partition leader, not the gateway's own compute.
  • Stamp priority once, at the edge, into the event header, and have every stage downstream read the same field. A pipeline that re-decides priority per stage sheds inconsistently and produces partial sessions, which is the failure the two-dimensional scheme exists to prevent.
  • Shed by a hash of the entity id, not at random. Dropping 30% of events uniformly corrupts every user's data slightly. Dropping 30% of users entirely leaves 70% exact and makes the gap explainable.

What does not transfer

The piggyback. A log has no synchronous response path back to the original producer, so there is nothing to attach an admission level to. The equivalents are slower and coarser: broker-side client quotas, an explicit control topic that producers poll, or a 429 from the ingest API that the client SDK must respect. The feedback loop runs in seconds rather than microseconds, which means the gateway needs a bounded accept queue so that overload is visible immediately rather than buffered into memory.

The second mismatch is granularity. DAGOR assigns priority per request, where one request is one user action. An ingest batch of 500 events can span many users, so priority has to be evaluated per record inside the batch, and a partial-batch rejection needs a response shape that tells the client which records to retry.

Where copying it would be a mistake

When the producers are first-party services you control, throttle at the source instead. Back-pressure applied where the data originates is strictly better than shedding at the gateway, because the producer can buffer, batch or degrade its own sampling rate with full context. Shedding exists for the case where you cannot reach the producer.

When every event is a financial fact, do not shed at all. Buffer at the edge with durable local storage and let freshness degrade instead of completeness. A pipeline that drops ledger events under load has converted an availability problem into a reconciliation problem, and reconciliation is the more expensive one.

Common weak answers

  • "Autoscale the ingest tier." Scaling takes minutes, the overload takes seconds, and the downstream broker partition count does not scale with the gateway.
  • "Add a bigger queue." A deeper queue converts rejection into latency, and the events at the back are stale by the time they are processed, so you pay to produce output nobody can use.