A live-streaming platform at Twitch scale load tests its chat-history service against two billion synthetic messages whose channel IDs were drawn uniformly at random. The test reports p99 of 40 ms at 5000 requests per second. Production at 3000 requests per second shows p99 of 800 ms with one shard at 95% CPU while the rest sit near 20%. Which explanation fits the evidence?
Show the full answer Hide the answer
The signal that settles it
The per-shard CPU spread is the whole diagnosis. One shard at 95% and the others at 20% is a statement about key distribution and nothing else. No amount of generator error, cache warmth, network topology or payload size produces a fivefold imbalance across shards; those all affect every shard roughly equally.
The mechanism
Drawing the partition key uniformly at random destroys the property that matters most for a partitioned store: the frequency distribution of the key. Live content is extremely skewed - a common planning shape is that the top 1% of channels carries around half of concurrent viewing, and during a major event a single channel can carry a significant share on its own. Uniform keys spread two billion messages perfectly across every shard, so the load test measured a system that cannot occur in production.
There is a second-order effect worth noticing because it explains why nobody caught it. With uniform keys over two billion messages the working set never fits cache, so the test's cache hit rate was worse than production's. The test was pessimistic on one axis and catastrophically optimistic on the other, and the pessimistic axis is the one that made the result look credible.
The fix
Sample the production key frequency distribution - counts per key only, no message content and no user identifiers, so nothing sensitive leaves production - and drive the generator from it. Preserve the head exactly: the top few hundred keys by frequency, at their real relative weights. Fit the tail with a power-law approximation, because the tail's exact shape does not change shard behaviour.
Then the real engineering question finally appears, and it is not a capacity question. A hot key is not fixed by a faster shard: it is fixed by splitting the hot channel's history across sub-partitions by time bucket, by serving its reads from replicas or a local cache, or by coalescing concurrent identical reads. The load test's job was to make that question visible before an event, and a uniform generator is precisely the thing that hides it.
Why the other options fail
- A fixed pool of virtual users under-reporting latency. This is real - it is coordinated omission, and it does make load tests optimistic. But it makes them optimistic uniformly. It cannot produce one hot shard, and the test's reported throughput of 5000 requests per second exceeds production's 3000, which is the opposite of what a saturated user pool looks like.
- A warm cache in the test. Plausible in general and contradicted here: two billion uniformly distributed keys give the test the worst possible cache behaviour, not the best.
- Cross-zone network latency in production. Cross-zone hops cost on the order of a millisecond, not 760. Wrong order of magnitude, and again it would be spread across shards.
- Smaller synthetic messages. Payload size shifts bandwidth and serialisation cost roughly linearly and hits every shard. It is worth fixing in the generator and it is not what is happening.
When this is the wrong thing to chase
If the partition key is a hash of an identifier with no semantic hot spot - a random UUID per record, never queried by a natural grouping - then uniform synthetic keys are correct and reproducing the production distribution is wasted work. The test is whether any real-world entity can concentrate traffic on one key. For channels, tenants, products, celebrities and postcodes the answer is yes, and it is the first thing to check about a generator.