Every day on YouTube, people upload 4 million videos and watch 5 billion videos. Handling this staggering traffic requires a vast fleet of servers. So when a new request comes in, where does it go? How do you balance load across so many servers for so many jobs?
I love this paper because it proposes a clever load balancing system that works at the largest of scales. Let’s imagine a YouTube service is distributed across 1000 servers. Whenever a new query arrives, we send it to one of those servers to process. How do we choose which one?
The obvious load balancing strategy is CPU-based, sending requests to whichever server has the lowest utilization. This works pretty well! But the problem with CPU-based load balancing is that CPU utilization is a lagging indicator that’s averaged over a measuring period and might not reflect the exact current state of a server. That means CPU-based load balancing can risk momentarily overloading servers and causing spikes in tail latency.
How do we balance load better? The idea is simple and elegant: if your goal is to minimize query latency, why not send each query to the server with the lowest latency? This works through probing. When a new query comes in, the load balancer sends lightweight probes to a number of target servers. Each probe measures two instantaneous signals: estimated latency of the server (based on latency of recent queries) and the number of requests in flight. Usually, queries are routed to the server with the lowest estimated latency. The exception is if the servers have very high numbers of requests in flight, indicating high incoming load–then, requests are routed to the servers with the fewest requests in flight.
What’s most impressive about this paper is that Prequal really works. The authors report that when they deployed it in production at YouTube, tail latency dropped by 2x, significantly reducing error and lag spikes for users. It’s not often that we see a systems paper produce results like that in the real world.