A search tier fans every query out to 60 shards and merges the replies. The merged p99 jumps from 40 ms to 240 ms, with a few requests at 440 ms. Shard CPU is flat, no individual shard's p99 moved, and the latency increments are almost exactly 200 ms. What failed?
Show the full answer Hide the answer
The trigger
Sixty shards receive a query at the same moment and answer within microseconds of each other, so their replies converge on one switch port as a synchronised burst. Top-of-rack egress buffers are shallow - a few hundred kilobytes for that port - and 60 shards returning 20 KB each is 1.2 MB arriving at once. The buffer overflows and packets are discarded. This is TCP incast, documented in data-centre networking research from the late 2000s and the motivation for much of the DCTCP line of work.
Why it propagates into a 200 ms step
A 20 KB reply is about 14 segments. When the dropped packet is near the end of a response, there are no following packets to generate duplicate acknowledgements, so fast retransmit cannot fire and the sender waits for a retransmission timeout instead. RFC 6298 recommends a one-second minimum; Linux clamps it at 200 ms, while datacentre round trips are under 250 microseconds. One drop therefore costs more than 800 round trips of an idle link.
The merge cannot finish until the slowest shard answers, so one loss among 60 replies sets the whole query's latency. That is why the distribution shows discrete 200 ms steps rather than a smear, and 440 ms when a second timeout follows with backoff.
Why detection lagged
Every application-level signal is innocent. Shards replied promptly, so their own latency is fine. CPU is flat because the loss is in the fabric. The signals that show it are not on the service dashboard: TCP retransmit counters on the aggregator's host and output-discard counters on the switch port.
The structural fix versus the tempting local fix
Tempting and wrong: retry adds more synchronised traffic to a congested port; raising the client timeout hides the steps; adding shards increases the fan-in degree and makes it worse.
Structural, strongest first:
- Reduce fan-in degree. A two-level merge tree turns 60 synchronised senders into eight groups of eight. This attacks the cause rather than the recovery.
- Jitter the fan-out by a few hundred microseconds so replies do not arrive as one burst. It costs a little median latency and removes the synchronisation.
- Keep tail loss recoverable without a timeout - SACK, tail loss probe and RACK all exist for this - and lower the RTO floor where the platform permits it.
- Keep queues short in the fabric with ECN-based congestion control, so the buffer is not already full when the burst lands.
Decision rule: if fan-in degree multiplied by response size exceeds the port's buffer, expect this. It is an arithmetic check anyone can do before the incident.
Common weak answers and the general lesson
- "Add capacity to the shards." They are not busy. Capacity is the answer to saturation and this is loss.
- "Set a tighter client timeout and retry the slow shard." A retry into a congested port is more of the cause, and a hedged request only helps if it takes a different path.
A latency distribution with discrete steps at a kernel timer value is a transport problem, not a capacity problem. Round numbers in a histogram - 200 ms, 1 s, 3 s - are timers, and timers point at retries and timeouts rather than at saturation.