Cloudflare operates at such a scale that even after many years of working here, it feels unreal. We have thousands of servers around the world with petabytes of RAM and millions of CPU cores, all running at capacity. As vast as these resources may seem, they are still finite, and when you need every service to run on every node, there is no room for inefficient space usage.
At this scale, small improvements have a huge impact, so even 1% optimizations at a time are worth paying attention to. And some improvements add up to a much greater result: in this article, we will look at how small changes in one algorithm significantly reduced the memory consumption of one of our Pingora-based services. This allowed us to free up over 100 TB of RAM globally, in addition to the 100 TB of memory the DNS team was able to save last month.
No waste
Maintaining a fair distribution of resources between teams is no easy task, especially in large organizations. One of the ways Cloudflare ensures this balance is through the tireless efforts of the wonderful Performance team.
This story begins with a ticket created by Ivan, who discovered: Excessive memory usage of pingora-ketama in Pingora Backend Router. It turned out that our internal load balancing service, Pingora Backend Router (yes, PBR), was using significantly more memory than expected, specifically in structures related to pingora-ketama, our open-source library for implementing consistent hashing.
To talk about how we handled this apparent memory over-allocation, we need to talk about what consistent hashing is, why we use it in PBR, and how it became so memory-hungry. Along the way, we will explore a bit of Rust and even touch on some math.
Consistent Hashing
Consistent hashing is a widely used method for distributing tasks across multiple servers in a way that does not require massive changes when adding or removing servers. Internally, we use it to route cached requests to servers by URL. This allows us to store only one copy of a file in each data center and provides a stable way to find the location of each file. We have mentioned this system before, but let's take a moment to break down in detail how and why this algorithm is used and how it works.
The key concept of consistent hashing is that while hash functions can take any input, their output is limited to a single unsigned integer (32-, 64-, or 128-bit integers depending on the hash function). This allows us to link tasks and servers together in a consistent manner. In most discussions of consistent hashing, the output space is represented as a continuous ring that loops from the maximum value back to zero. This description allows for beautiful visualizations, but it can also make a simple concept of integer ranges more complex than necessary. For our discussion, we will represent the 32-bit output of our hash function as a number line.
Now, suppose we have a set of servers A, B, and C and a set of tasks t-z. We can map each of them to the number line based on the hash of their representative values, such as IP addresses for servers and cache keys for tasks.
Assigning tasks to servers now boils down to finding the first server to the left of each task. We can represent this visually by coloring the hash area that will be associated with each server. Note that the range covered by server C loops back to the beginning, hence the idea that hashes exist as a ring.
That's it. At a basic level, consistent hashing is very simple, but it soon becomes obvious that there is room for improvement. Note that the range covered by server A in our example is significantly larger than that of B or C. This is a problem because the share of requests handled by a server will be proportional to the size of its range on the number line. Ideally, we would like to ensure that every server has the same size, but since hashes are essentially random numbers, we have to talk about region sizes in terms of statistics. 😨
Math and its consequences
First: don't panic. I promise I'm not going to trick you, and we will stay within the safe confines of an introductory probability lesson. When we talk about statistical distributions, there are two important factors that help us estimate uncertainty in a useful way: the expected value and the standard deviation. To put it (too) simply, the expected value gives us the point around which measurements based on the distribution will be concentrated, and the standard deviation shows how close most measurements are to that central point.
For consistent hashing, we can calculate these factors for the relative size of the range associated with one of N servers. (Details on where this formula comes from will be later).
$$m \begin{align*} \text{Exp} &= \frac{1}{N} \\ \text{SD} &= \frac{1}{N}\sqrt{\frac{N-1}{N+1}} \end{align*} m$$
Speaking of specific numbers, suppose we have 100 servers. The formulas above give:
$$m \text{Exp}=1/100 = 1\% \\ \text{SD}= \frac{1}{100}\sqrt{\frac{100-1}{100+1}} \approx 0.99\% m$$
This tells us that we can expect the range handled by each server to be centered around 0.99% of the total, and most lengths will fall within 1% of the expected value. This sounds okay until we realize that this is 0.99% of the total length. We need to scale the standard deviation by the expected value to see how large the error is as a fraction of the target size. This value is called the coefficient of variation.
$$m \text{CV} = \frac{\text{SD}}{\text{Exp}} = \sqrt{\frac{N-1}{N+1}} m$$
At $m N=100, \text{CV} \approx 99\% m$ — this means that some servers are likely to be working 99% harder than they should (handling twice as many requests), while others may be doing almost nothing! Now that we have a way to predict load uniformity on servers when using consistent hashing, we can start making improvements.
What if we add hashes?
The simplicity of consistent hashing is a double-edged sword. It is easy to understand and implement because everything is converted into easily comparable hashes on the same number line, but any improvements to the system must also be tied to this number line. This means that the solution to any consistent hashing problem can only be new hashes. It is less like a golden hammer (a tool that makes every problem look like a nail) and more like a golden nail, as it turns every tool into a hammer.
To solve the problem of unbalanced workloads, we can add multiple hashes to represent each server instead of just one. We will soon move on to the mathematical justification for this process, but it should be intuitively clear that while each individual range has a large standard deviation, adding many such ranges leads to their total size leveling out. If we take our three-server example from the diagrams above and add two more random hashes for each server, we will see that this helps to balance the workload of each server.
This is, admittedly, an artificial example. The random nature of the system means there are no guarantees as to how much the situation will improve by adding 2 extra hashes per server, but it should be intuitively clear that combining more of these hash segments together leads to a more uniform distribution. Each segment in total has a chance to compensate for another. Perhaps one is too short; perhaps one is too long. This is essentially what the law of large numbers tells us... The obvious problem is that it only works for large numbers. In NGINX, the base number of hashes per server is hardcoded to 160, and Pingora uses the same default value. I will spare you the math for now, but if we go back to our 100-server example and use 160 points per server instead of one, the coefficient of variation (which can be thought of as the margin of error) will drop from about 99% to about 8%, which is a significant improvement.
What if we add more hashes?
Above, we saw that increasing the number of hashes per server by a constant amount allows for improving the uniformity of workload distribution across servers, but what if we don't want to distribute work evenly? In Cloudflare's case, we have servers with more disk space than others, so it would be better if the number of requests allocated to a server were proportional to the size of its disk. One way to achieve this is the ketama algorithm. The name sounds a bit funny because the algorithm is named after the library where it was first implemented, and the library was named... well, you can google that 😶🌫️.
The entire algorithm boils down to the following: for any two servers, $m S_1m$ and $mS_2m$, if we want the requests handled by $mS_1m$ to be $mw\timesm$ of the requests handled by $mS_2m$, the number of hashes associated with $mS_1m$ must be equal to $mH_1 = w\times H_2m$. This allows us to set a "weight" for each server, which scales the number of hashes associated with that server. Unfortunately, this does not replace the constant scale factor we added in the previous section. Such scaling is necessary to set the minimum margin of error that will manifest on servers with the lowest weights.
For us, since we want the workload to scale based on storage, we can use disk space as the weight, which is what the Pingora team has been doing for many years. In other parts of the company where workloads are more compute-intensive, weights can be based on the number of CPUs or GPUs.
What if we add even more hashes???
The final problem we need to solve is that, until now, we have been working under the assumption that any server can handle any request, but in practice, this is not the case. Things like compliance requirements or enabled caching features mean that only a subset of servers can handle a specific request. Unfortunately, unlike the previous case, we cannot solve this problem by adding more hashes to the same ring. We have to add entirely new rings, and not only that—every combination of features potentially needs its own specific ring!
Duplication based on combinations is a classic recipe for exponential explosion. In our case, we have a handful of different features, leading to $m2^\text{handful} = \text{dozens}m$ of separate consistent hashing rings. So, as you probably already guessed, the "excessive memory usage" (in some cases 6 GB) discovered by Ivan was driven by the huge number of hashes needed to provide all the functionality we require, which must be stored in RAM.
Storage Improvements
One major improvement was proposed by Zaidoon, who had a realization about storing hashes in PBR. This structure looks like this:
In memory, it is represented as eight bytes, where four go to the hash (which cannot be avoided) and four to an index pointing to the server, which is stored in another array. Zaidoon's idea was that a 32-bit integer for this index is redundant, as PBR is unlikely to ever need to coordinate more than $m2^{16} \approx 65\text{k} m$ servers simultaneously, so a 16-bit integer will suffice. Thus, we can replace the structure above with this one:
Unfortunately, Rust does not make this simple. Changing the index size as we did above does not reduce the amount of memory occupied at all. This is because Rust has alignment rules that require the size of a structure in memory to be a multiple of its largest (or "most aligned") field. In this case, the hash is the largest with four bytes, so when stored in memory, Point must have a size of $mN \times 4m$, meaning the minimum size is eight bytes.
Fortunately, there are well-known ways to get around this. You (that is, me) might be tempted to use #[repr(packed)], but this is controversial for good reasons. A safer, though less readable, solution is to store the hash and index as a raw array of bytes and access them using getters. Both methods compile to the same thing.
This simple (though verbose) change reduces the amount of memory used for consistent hashing by a full 25%! To achieve more, we will have to return to the math, so hold on tight; this is the home stretch.
What if we try fewer hashes?
You might have noticed that we provided the standard deviation formula for the case where there is only one hash per server. Deriving the formula for the case where there are $m k m$ hashes per server is not straightforward, and most sources provide only an approximation or an asymptotic limit, but not us. I may not be a statistician, but I grew up with a math teacher (hi, Mom!), and I wanted to know the real value. The full derivation is provided in an additional article, but here is the result.
$$m \text{Exp}_k = \frac{1}{N}, \text{SD}_k=\sqrt{\frac{(k+1)}{N(kN+1)}-\frac{1}{N^2}} m$$
To see how increasing the number of hashes improves accuracy, we need to look at the coefficient of variation again.
$$m \text{CV}_k=\frac{\text{SD}_k}{\text{Exp}_k}=\sqrt{\frac{N-1}{(N*k+1)}} m$$
Plotting $m\text{CV}_km$ reveals a potential problem with the "just add more hashes" mentality (aside from RAM overhead).
You can see that each step of error reduction requires (almost) an order of magnitude larger increase in the number of hashes per server, so adding new hashes yields smaller and smaller improvements. Recall that we use a base of 160 hashes, scaled by the server's storage size. To simplify the math, let's assume the weight factor $m{m_w}m$ for a server is 625, then we get $m{k = 160\times625 = 100{,}000}m$. From the graph above, it is clear that the last 90,000 hashes we added provide a negligible 0.7% reduction in error. Unfortunately, it only gets worse from there.
Predictions based on my beautiful math only work when we consider hashes as a continuous ring, but in practice, we use 32-bit numbers for hashes, which can potentially lead to collisions, and the probability of collisions increases surprisingly quickly as the number of hashes increases (see the birthday paradox). Collisions matter because, in an ideal case, each hash contributes to the volume and distribution of requests handled by the associated server, but a collision means that some contributions are randomly discarded, introducing unpredictable error. If we compare the simulation results with 32-bit hashes to the predicted error rate, we see that for data centers with 2048 servers, the error rate increases in the range of 10,000 to 100,000 hashes per server.
Ultimately, while this realization doesn't seem the most pleasant, it is great news for our plan to free up RAM! Now that we have the mathematical justification, we have determined that we can reduce the number of hashes generated for each server by 90% without causing any noticeable error, and that is exactly what we decided to do.
Migration without overloading origin servers
There was another problem: changing the hash ring changes the destination of some cached requests. Even if the new ring is better, switching the entire network at once would effectively invalidate almost all cached content. This would turn memory optimization into a catastrophic traffic spike on the origin servers.
Therefore, we did not do this as a single global switch. For a while, PBR kept both versions of the load balancer for cached content in memory: the old ketama ring and the new, smaller one. Each request used our usual migration framework to determine which ring the backend should choose. This meant that the deployment decision was stable for each request hash and also gave us a clean rollback path. If something went wrong, we could route new requests back through the old ring without redeploying PBR.
Then we performed the migration in stages. We started with small validation locations, moved to gradually increasing groups of data centers, and only then continued for the rest of the world.
An important part was that we independently controlled two dimensions: what fraction of traffic uses the new ring and where that traffic is allowed to move. A simple deployment with a global percentage would have spread cache wear everywhere at once. Deploying within data centers allowed us to keep a small blast radius and made it significantly easier to understand whether the change was actually safe.
During the migration, we monitored backend selection traces, ring version counters, PBR connection errors, process memory, startup time, cache behavior, and origin traffic. Once the migration reached 100%, we removed the temporary path with the old ring, and voila!
The graph above shows a comparison of the memory used by PBR during the week of the change with data from several weeks prior, as well as the result of subtracting one from the other. The sharp drop is the day the PBR version with the large (now unused) hash rings was decommissioned forever. Looking at the difference, we get an impressive result: our changes reduced memory usage by 100 TB!
Try it yourself
All the changes we described in this post are now available in the pingora-ketama crate as a (not yet announced) cargo feature. The v2 ring has a compact storage format, a faster sorting method, and the ability to scale the base number of hashes per node. Our priority when making these changes was stability and control, so the v1 ring is identical to what pingora ketama has always used, and the library allows you to run them both simultaneously and decide for each request individually which one to use and when.
Beyond testing our hashing changes, I would like you to take away inspiration to examine your own systems for what "simple" or "obvious" solutions hide potential benefits if you are willing to crunch the numbers. You might not be able to solve all your problems with Rust, but math is universal.










