Adding a node changes capacity. It does not prove the load is evenly placed.
At 14:05, shard 3 has a growing write queue and elevated storage latency. The cluster dashboard shows six nodes around 65% CPU, while shard 3 reaches 94%. The service adds a seventh node; the new node sits mostly idle because the partition map still sends the same hot customer keys to shard 3. A rebalance starts, but the queue remains high while migration competes for disk and network bandwidth.
Those observations support a placement or workload-skew investigation. They do not yet tell us whether the hash function is uneven, whether a few keys dominate writes, whether one node has a slower disk, or whether requests are routed with stale metadata. The first useful artifact is a per-shard view of key count, request rate, bytes, and resource work over the same time window.
- Observed
- Shard 3 is at 94% CPU with increasing write wait.
- Changed
- A seventh node joined; its CPU remains low.
- Unknown
- Are key placement and request work both evenly distributed?
- Risk
- Moving partitions may add load before it removes load.
Equal key counts can hide very unequal work.
For each node i, record a named load measure Lᵢ over one interval:
requests/s, bytes/s, CPU-seconds/s, queue depth, or a workload-specific cost estimate. Keep
units and boundaries visible. A row count is useful for storage balance; it is not a substitute
for write amplification or CPU time if those are the scarce resources.
A simple first view is the largest-to-mean ratio: max(Lᵢ) / mean(Lᵢ). The
coefficient of variation, CV = standard deviation(Lᵢ) / mean(Lᵢ), summarizes
spread relative to the average. Neither metric tells you whether the mean is safe, whether
an outlier is a hot key, or whether demand is correlated in time. Show the node values and
a time series alongside any single summary.
Modulo makes placement cheap; changing the node count changes the answer for most keys.
A common rule hashes a key to an integer and selects hash(key) % N. If hash
values are uniform and keys have equal weight, this spreads many keys approximately evenly
across N nodes. With K independent uniform keys, the count on one node
has mean K/N and standard deviation about √(K(1/N)(1 − 1/N)). Relative
noise shrinks as key count grows, but it never promises equal counts.
The bigger operational cost appears when N changes. A key that formerly
mapped to hash % N may now map elsewhere under hash % (N + 1). Under ideal
uniform hashing, the expected fraction that stays on its old node when adding one node is
roughly 1 / (N + 1); so about N / (N + 1) keys move. At six nodes, that is
roughly 86% moved. This may be acceptable for a stateless routing decision, but expensive when
the mapping owns cached or persistent data.
That estimate assumes a uniform independent hash and keys treated equally. Real implementations must also define hash compatibility, integer width, signedness, and how membership is discovered.
Consistent hashing limits movement by preserving most existing key owners.
A consistent-hash ring places each node at one or more positions on a circular hash space.
A key is hashed to a position, then assigned to the next node position clockwise. When a
node joins, only keys in the interval it takes over move; when a node leaves, only its
keys move to their next owners. With equal nodes and ideal uniform points, adding one node
to N existing nodes moves about 1 / (N + 1) of keys on average, instead
of almost all keys.
One position per node can create large accidental gaps and uneven ownership. Virtual nodes assign several ring positions to each physical node, making the aggregate ownership more even. More positions reduce placement variance in expectation; they do not guarantee balanced request cost. They also increase ring metadata, hashing work, and reconfiguration complexity. Unequal capacity can be represented by assigning more virtual positions to stronger nodes, but the weights need measurement and careful rollout.
Use a stable hash and a stable encoding.
Walk to the next node position, wrapping once.
Most keys keep their current owner.
Power-of-two choices uses current load to avoid the most crowded of two candidates.
For a request, hash its identity twice to obtain two candidate nodes, inspect a load signal, and choose the less-loaded candidate. With many independent requests, this simple choice can reduce maximum load sharply compared with assigning every request to one random node. The improvement comes from adapting the decision to current occupancy, not from a more even hash by itself.
The method needs a timely, comparable load signal and a tie rule. Concurrent routers may see stale loads and herd toward the same apparently quiet node. The load signal might be active requests, queue depth, or a measured cost estimate; choose one that matches the bottleneck and account for measurement delay. The algorithm spreads independent routing decisions, but it can break affinity if the same key must always reach the same owner.
A single hot key is a special case. Sticky modulo or ring placement sends every operation for that key to one owner; hashing the key more carefully cannot split the key's work. Replicate read-only data, split a hot partition into subkeys, cache, or route requests dynamically only if the consistency and ordering model allows it. Measure the resulting coordination and write costs.
Compare placement spread and the keys that move when one node joins.
This deterministic lab hashes equally weighted sample keys with one illustrative 32-bit hash. Modulo and ring rows model sticky key placement; two choices models sequentially assigning independent requests to the currently less-loaded of two candidates. It is a comparison of mechanisms, not a production benchmark.
| Node | Modulo | Ring | Two choices |
|---|---|---|---|
| 0 | 103 | 18 | 102 |
| 1 | 96 | 167 | 99 |
| 2 | 95 | 25 | 99 |
| 3 | 102 | 51 | 99 |
| 4 | 102 | 288 | 100 |
| 5 | 102 | 51 | 101 |
The theoretical moved share under ideal uniform hashing is about 14.3% for a ring and 85.7% for modulo. This finite sample may differ. More virtual positions can smooth ownership, while request skew and hot keys remain outside this equal-weight model.
Keep stable placement and load-aware routing visibly distinct.
These small examples share a 32-bit FNV-1a hash over UTF-8 bytes for cross-language parity. They show modulo ownership, a consistent ring with virtual nodes, and a power-of-two choice using caller-supplied current loads. Production systems need a deliberate hash choice, compatible serialization, operational membership updates, and concurrency-safe load measurements.
The ring minimizes remapping; the two-choice helper is dynamic and does not preserve key affinity.
package mathpractice
import (
"errors"
"math"
"sort"
)
type RingPoint struct {
Position uint32
Node int
}
// Hash32 is a small deterministic 32-bit FNV-1a over UTF-8 bytes, not a security hash.
func Hash32(value string) uint32 {
hash := uint32(2166136261)
for i := 0; i < len(value); i++ {
hash ^= uint32(value[i])
hash *= 16777619
}
return hash
}
// ModuloOwner uses simple placement; changing nodeCount can move most keys.
func ModuloOwner(key string, nodeCount int) (int, error) {
if nodeCount < 1 {
return 0, errors.New("node count must be positive")
}
return int(Hash32(key) % uint32(nodeCount)), nil
}
// BuildRing places each node at virtualNodes points on the ring.
func BuildRing(nodeCount, virtualNodes int) ([]RingPoint, error) {
if nodeCount < 1 {
return nil, errors.New("node count must be positive")
}
if virtualNodes < 1 {
return nil, errors.New("virtual node count must be positive")
}
ring := make([]RingPoint, 0, nodeCount*virtualNodes)
for node := 0; node < nodeCount; node++ {
for replica := 0; replica < virtualNodes; replica++ {
ring = append(ring, RingPoint{
Position: Hash32("node:" + itoa(node) + ":replica:" + itoa(replica)),
Node: node,
})
}
}
sort.Slice(ring, func(i, j int) bool {
if ring[i].Position == ring[j].Position {
return ring[i].Node < ring[j].Node
}
return ring[i].Position < ring[j].Position
})
return ring, nil
}
// ConsistentOwner selects the first clockwise point, wrapping to the first point.
func ConsistentOwner(key string, ring []RingPoint) (int, error) {
if len(ring) == 0 {
return 0, errors.New("ring must not be empty")
}
position := Hash32(key)
index := sort.Search(len(ring), func(i int) bool { return ring[i].Position >= position })
if index == len(ring) {
index = 0
}
return ring[index].Node, nil
}
// PowerOfTwoChoice routes one movable request using the current observed loads.
func PowerOfTwoChoice(requestID string, loads []float64) (int, error) {
if len(loads) < 2 {
return 0, errors.New("power-of-two choice needs at least two nodes")
}
for _, load := range loads {
if math.IsNaN(load) || math.IsInf(load, 0) || load < 0 {
return 0, errors.New("loads must be finite and non-negative")
}
}
first := int(Hash32("choice-a:"+requestID) % uint32(len(loads)))
second := int(Hash32("choice-b:"+requestID) % uint32(len(loads)))
if second == first {
second = (second + 1) % len(loads)
}
if loads[first] <= loads[second] {
return first, nil
}
return second, nil
}
// itoa keeps the sample dependency-free; values here are non-negative indexes.
func itoa(value int) string {
if value == 0 {
return "0"
}
var digits [20]byte
i := len(digits)
for value > 0 {
i--
digits[i] = byte('0' + value%10)
value /= 10
}
return string(digits[i:])
}
Change routing only after the measurements show which imbalance you have.
Before a rebalance, capture a baseline over a representative interval: per-node key counts, request rates, bytes, resource cost, queue wait, and tail latency. Identify the hottest keys and whether traffic for them is read-heavy, write-heavy, or bursty. Estimate the data and cache movement and the extra migration load. Then compare a small canary or replay of the same trace under the candidate policy, with rollback conditions tied to customer-visible service measures.
Choose modulo when membership is stable and broad remapping is acceptable. Consider a consistent ring when key ownership must survive membership changes with limited movement, and tune virtual nodes against observed placement variance and operational cost. Consider two choices when work is independent and movable, with a load signal whose staleness and coordination costs are understood. If one identity dominates, address its workload or split/replicate it under explicit consistency rules; no distribution algorithm can make an indivisible hot key cease to be hot.