Consistent Hashing Limits Key Movement During Topology Changes
Consistent Hashing Limits Key Movement During Topology Changes A distributed cache or partitioned service needs a rule that maps each key to a node. A simple rule such as hash(key) % N is attractive while the node count stays fixed. The trouble appears when N changes. Moving from four nodes to five changes the divisor for every key. Most remainders change, so a routine capacity adjustment can remap a large share of the dataset at once. For a cache, that can trigger a wave of misses. For stateful storage, it can create a large migration job.