已發布·持續改進
一致性雜湊 指南 · 6/6
本章目前僅提供英文版。
Consistent hashing began with a 1997 paper by Karger and colleagues at MIT on distributing web caches, and that research fed into the content delivery network Akamai. Today it shows up in distributed caches, key-value stores and load balancers. Not every system that "shards" uses it, though, so check what a given system actually does.
Memcached servers know nothing about each other; the client decides which server holds a key. Ketama, released by last.fm in 2007, gives each server many points on a ring, and libmemcached along with several other clients and proxies such as twemproxy offer ketama-compatible distribution. Adding a cache server leaves most keys where they were, so the hit rate does not collapse.
Redis Cluster does not use a consistent hashing ring. It maps each key to one of 16,384 hash slots with CRC16(key) mod 16384 and assigns slots to nodes explicitly. When a node is added, an operator moves some slots, so the amount of movement is similarly small, but the mechanism is different.
Amazon's 2007 Dynamo paper combined a consistent hashing ring, virtual nodes and a preference list that replicates each key to the next N nodes clockwise, and it was hugely influential. The paper also reports that splitting the ring into fixed, equal-sized partitions assigned to nodes worked better for balance and operations. Note that the AWS DynamoDB service was inspired by the paper but is a separate system with a different internal design.
num_tokens), and the default partitioner uses Murmur3. Replication walks the ring for the next nodes; NetworkTopologyStrategy also takes data centers and racks into account.Sending the same user or URL to the same backend lets caches and sessions be reused. Nginx switches its hash directive to ketama consistent hashing when you add consistent.
upstream cache_backend {
hash $request_uri consistent;
server 10.0.0.11:8080;
server 10.0.0.12:8080;
server 10.0.0.13:8080;
}HAProxy does the same with hash-type consistent, and hash-balance-factor turns on consistent hashing with bounded loads, which caps any server at a multiple of the average load. A value of 150 means 1.5 times the average.
backend cache_servers
balance uri
hash-type consistent
hash-balance-factor 150
server c1 10.0.0.11:8080 check
server c2 10.0.0.12:8080 checkEnvoy offers RING_HASH and MAGLEV load-balancing policies. Maglev, from Google's 2016 software load balancer paper, replaces the ring with a fixed-size lookup table for O(1) lookups.
Hot keys can overload a single node. The bounded-loads scheme by Mirrokni and colleagues caps every node at c times the average load and, when the owner is full, moves on clockwise to the next node. Here is a simple version on top of the implementation chapter's HashRing.
import bisect, math
def get_bounded(ring, key, load, total, c=1.25):
nodes = set(ring.owner.values())
cap = math.ceil(c * (total + 1) / len(nodes)) # per-node capacity
pts = ring.points
start = bisect.bisect_left(pts, hash64(key))
for step in range(len(pts)):
node = ring.owner[pts[(start + step) % len(pts)]]
if load.get(node, 0) < cap:
load[node] = load.get(node, 0) + 1
return node
raise RuntimeError("no capacity")hash(), moves every key on every restart.node#i) effectively rebuilds the whole ring.def read(key, new_ring, old_ring, stores):
value = stores[new_ring.get(key)].get(key)
if value is None: # not migrated yet
value = stores[old_ring.get(key)].get(key)
return value
0 則留言
登入 · 登入後即可留言。
來留下第一則留言吧。