TT Lab
Get started
Learn Learning paths Courses

Envoy Internals

Why Session Stickiness Sometimes Breaks

Continue in TT Lab

In one line

Once the cluster has narrowed down the "places worth sending to", picking one of them is load balancing. The ways of picking split by how much state they remember — remember nothing and it is random, remember whose turn it is and it is round robin, remember the key of the request and it is a hash ring. Whatever you pick comes with a cost.

Why this was needed

With three servers, it seems enough to split requests into thirds. In practice that is true most of the time. The trouble comes when two requirements arrive.

One is "please always send this user to the same server". If the server holds a local cache or keeps sessions in memory, this requirement decides both performance and correctness. The other is "when I remove one server, it is a problem if everything gets reshuffled". If servers holding caches all change places at once, then the moment you remove one, the overall cache hit rate falls to the floor and the backend takes that load.

The thing that solves these two requirements together is consistent hashing, and in Envoy that is RING_HASH and MAGLEV.

How it works

lb_policy How it picks Cost
ROUND_ROBIN One at a time, going around the list Each worker remembers its turn separately
LEAST_REQUEST Picks two at random and takes the one with fewer requests in progress Advantageous when request costs are uneven
RANDOM Random every time Nothing to remember, so it is cheap. It skews with a small sample
RING_HASH Hashes the key and takes the nearest endpoint on the ring The cost of building the ring. If the ring is sparse, the distribution skews
MAGLEV The same property more cheaply with a fixed-size lookup table The table size is fixed, so there is a constraint when there are very many endpoints

You can give each endpoint a load_balancing_weight. A weight of 2 does not mean "one more time out of every two" but that it takes twice the share of the total. If the weights of three are 2, 1 and 1, the sum is 4, so out of 12 requests it is 6, 3 and 3.

Hash-based methods need a key. What chooses it is the route's hash_policy, and you can pick from a header, a cookie, a query parameter or the source IP. One thing often missed here — a request whose key could not be extracted has no hash, so it just goes to a random place. The first request without a cookie and an internal call that left out the header are examples. This is usually what is behind "most of them are pinned, but sometimes one goes elsewhere".

The properties of the ring are worth knowing too. If you simply use "the hash value modulo the number of servers", almost every key moves the moment the number of servers changes. With the ring method, only the keys that were attached to the vanished server move, and the rest stay as they are. In exchange, if the ring is sparse (if minimum_ring_size is small), endpoints are not spread evenly over the ring, so the distribution itself skews.

What it looks like in the field

"I sent 9 requests and did not get 3, 3, 3." Each Envoy worker thread remembers its own round robin turn. The default concurrency is the number of cores, so if you count with few requests, they scatter across workers and the numbers differ every time. An experiment that counts the distribution uses --concurrency 1. In production, do not count; look at the statistics.

"I turned on session affinity but it sometimes comes loose." First check whether requests with no key are mixed in. And at the moment endpoints come and go, it is normal for some to move — consistent hashing does not mean "nobody moves" but "the fewest possible move".

When only one particular server runs hot. This happens when the distribution of the hash keys is skewed. If you used the tenant ID as the key and one tenant is half of the traffic, the server that tenant is attached to takes half. In that case, split the key finer (per user) or take that tenant out separately.

The criteria for choosing a method. Summed up, it comes down to two questions. First, must the same request go to the same place? If a local cache or sessions are involved, yes, and then you go hash-based. Second, are request costs uneven? If some requests take 5 milliseconds and some take 2 seconds, a method that splits in turn keeps piling more work onto the server that got a slow request. In that case LEAST_REQUEST, which looks at the number of requests in progress, is better. If neither applies, leaving it at round robin is the easiest to predict, and when the scale gets very large, stateless random actually becomes cheaper.

Official documentation: Load balancing overview · Cluster configuration · HTTP route components

What you will do in the next lab

You keep the same set of upstreams, change only the method, and count the numbers. You check in turn that round robin splits exactly, that weights change the shares, and that random skews on a small sample, and with the hash ring you count for yourself that the same user sticks to one place and how many users change places when one server is removed. Finally you change the key from a header to a query parameter and see that requests without the key are not pinned.