Consistent Hashing for Queue Shards
When one broker instance cannot carry a queue, the queue is split across shards, and the function that maps a job's key to a shard decides how painful it will be to add the next shard. This guide replaces modulo hashing with consistent hashing for queue shards, as part of Queue Partitioning Strategies in Queue Fundamentals & Architecture.
Problem Statement
A notification pipeline shards its per-tenant job queues across four Redis instances with shard = hash(tenant_id) % 4. Tenants must stay on one shard because each tenant's jobs are processed in order by a single consumer. Growth requires a fifth shard. With modulo hashing, changing 4 to 5 moves 80% of tenants to a different shard — every one of them has jobs in flight on the old shard and new jobs landing on the new one, breaking per-tenant ordering and requiring a coordinated migration of most of the data. You want to add shards moving only about a fifth of tenants, a way to move those tenants without reordering their jobs, and even load across shards.
Prerequisites
- A stable routing key per job (here
tenant_id) that all producers and consumers use. - A shared, versioned shard map that producers and consumers can read (config file, key-value store, or service).
- The ability to pause enqueue briefly for a single tenant (a per-tenant flag) during migration.
- Per-shard metrics for depth and throughput.
Step 1 — See How Much Modulo Hashing Moves
With hash(key) % N, changing N changes the result for almost every key: when going from N to N+1 shards, only about 1/(N+1) of keys keep their shard.
import hashlib
def h(key: str) -> int:
return int.from_bytes(hashlib.blake2b(key.encode(), digest_size=8).digest(), "big")
tenants = [f"t-{i}" for i in range(100_000)]
moved = sum(1 for t in tenants if h(t) % 4 != h(t) % 5)
print(f"{moved / len(tenants):.0%} of tenants change shard going 4 -> 5") # ~80%
Every moved tenant is a potential ordering break and a data migration. The ideal when adding one shard to N is to move only 1/(N+1) of keys — the ones the new shard should own — and leave everyone else alone. That is exactly what consistent hashing achieves.
Step 2 — Build a Hash Ring with Virtual Nodes
A hash ring places each shard at many points ("virtual nodes") on a circle of hash values. A key belongs to the first shard point clockwise from the key's hash. Adding a shard inserts new points; only keys between a new point and its predecessor move, and they all move to the new shard.
import bisect
class HashRing:
def __init__(self, shards: dict[str, int], vnodes: int = 256):
"""shards: name -> weight (relative capacity)."""
self.points: list[int] = []
self.owner: dict[int, str] = {}
for name, weight in shards.items():
for i in range(vnodes * weight):
p = h(f"{name}#{i}")
self.points.append(p)
self.owner[p] = name
self.points.sort()
def shard_for(self, key: str) -> str:
i = bisect.bisect(self.points, h(key)) % len(self.points)
return self.owner[self.points[i]]
ring4 = HashRing({"redis-0": 1, "redis-1": 1, "redis-2": 1, "redis-3": 1})
ring5 = HashRing({"redis-0": 1, "redis-1": 1, "redis-2": 1, "redis-3": 1, "redis-4": 1})
moved = [t for t in tenants if ring4.shard_for(t) != ring5.shard_for(t)]
print(f"{len(moved) / len(tenants):.1%} moved; all to redis-4: "
f"{all(ring5.shard_for(t) == 'redis-4' for t in moved)}")
# 20.1% moved; all to redis-4: True
Virtual nodes are what make the load even: with one point per shard, arc lengths vary wildly and one shard may own twice another's keys. With 256 points per shard, each shard's share is within a few percent of equal. Weights let a larger instance own proportionally more.
Step 3 — Consider Rendezvous Hashing for Small Shard Counts
Rendezvous (highest-random-weight) hashing achieves the same minimal movement without a ring: for each key, score every shard and pick the highest. It is simpler, needs no sorted structure, and is exact rather than statistical — at the cost of O(shards) work per lookup, which is trivial for tens of shards.
def rendezvous(key: str, shards: list[str]) -> str:
return max(shards, key=lambda s: h(f"{s}:{key}"))
shards4 = ["redis-0", "redis-1", "redis-2", "redis-3"]
shards5 = shards4 + ["redis-4"]
moved = sum(1 for t in tenants if rendezvous(t, shards4) != rendezvous(t, shards5))
print(f"{moved / len(tenants):.1%}") # ~20%, all to redis-4
For queue shard counts (usually under 100), rendezvous hashing is often the better choice: fewer moving parts and perfectly even distribution in expectation.
Step 4 — Publish the Shard Map as Versioned Configuration
Producers and consumers must agree on the shard set at all times. Put the map in one place with a version, and make both sides read it.
# shard-map.yaml (stored in a config service; clients cache and watch for changes)
version: 7
algorithm: rendezvous
shards:
- { name: redis-0, url: "redis://redis-0:6379", weight: 1 }
- { name: redis-1, url: "redis://redis-1:6379", weight: 1 }
- { name: redis-2, url: "redis://redis-2:6379", weight: 1 }
- { name: redis-3, url: "redis://redis-3:6379", weight: 1 }
- { name: redis-4, url: "redis://redis-4:6379", weight: 1, state: warming } # not yet routable
overrides: {} # tenant_id -> shard, for pinned or migrating tenants
The state: warming flag lets the new shard be deployed and monitored before any traffic is routed to it; overrides pins specific tenants during migration (Step 5) or permanently (a huge tenant on a dedicated shard, as in partitioning queues by tenant ID).
Step 5 — Move Tenants Without Reordering Their Jobs
For each tenant that will move, jobs already queued on the old shard must finish before new jobs start on the new shard. Move tenants one at a time (or in small batches) with a drain-then-switch sequence.
def migrate_tenant(tenant: str, old: str, new: str) -> None:
# 1. Pin the tenant to its current shard while preparing
shard_map.set_override(tenant, old)
# 2. Pause enqueue for this tenant: producers buffer or delay briefly
tenant_flags.set(tenant, "enqueue_paused", ttl=600)
# 3. Wait until the tenant's queue on the old shard is empty and nothing is in flight
wait_until(lambda: queue_depth(old, tenant) == 0 and in_flight(old, tenant) == 0, timeout=300)
# 4. Switch: override to the new shard (equal to the new computed owner)
shard_map.set_override(tenant, new)
# 5. Resume enqueue; new jobs land on the new shard, in order
tenant_flags.clear(tenant, "enqueue_paused")
for t in tenants_moving_to("redis-4"):
migrate_tenant(t, ring_old.shard_for(t), "redis-4")
# Finally: publish map version 8 with redis-4 active and clear overrides that equal the computed owner
Pausing one tenant for a few seconds is invisible to users when producers retry or buffer; pausing everyone for a migration is not. With consistent hashing, only a fifth of tenants go through this; with modulo hashing, four-fifths would.
Verification
def test_distribution_is_even():
counts = Counter(rendezvous(t, shards5) for t in tenants)
assert max(counts.values()) / min(counts.values()) < 1.05 # within 5%
def test_adding_shard_moves_only_its_share():
moved = [t for t in tenants if rendezvous(t, shards4) != rendezvous(t, shards5)]
assert abs(len(moved) / len(tenants) - 0.2) < 0.01
assert all(rendezvous(t, shards5) == "redis-4" for t in moved)
In production, compare per-shard enqueue rates after the migration: they should converge to within a few percent unless a few very large tenants dominate, in which case pin those tenants explicitly.
Gotchas & Edge Cases
Hot keys defeat even hashing. Consistent hashing spreads keys evenly, not load. One tenant with 30% of traffic makes its shard hot regardless of algorithm. Pin heavy tenants to dedicated shards or split their key.
Inconsistent maps between clients. A producer with map version 7 and a consumer with version 8 disagree on where a tenant lives. Roll out map changes with overrides first (which both versions honour), then the algorithmic change.
Hash function choice. Use a stable, well-distributed hash (BLAKE2, xxHash, murmur3). Language built-ins like Python's hash() are randomised per process.
Removing a shard. The reverse operation moves only the removed shard's keys, spread across the remaining shards. Drain it the same way before decommissioning.
FAQ
Does Redis Cluster already do this? Redis Cluster uses 16,384 hash slots mapped to nodes, and moves slots when resharding — a similar minimal-movement idea. It does not know about your ordering requirements, so jobs in migrating slots can still interleave; the drain-then-switch step is still needed for ordered work.
What if producers are written in several languages? Every producer and consumer must compute the same shard for the same key, so the hash function, its input encoding, and the algorithm must match exactly across languages. Pick a hash with identical implementations everywhere (xxHash64 or BLAKE2b with a fixed digest size), hash the UTF-8 bytes of the key, and ship a shared test vector — a list of keys with their expected shard for a given shard set — that each language's client must pass in CI.
Should workers also use the ring? Workers usually do not need it: each worker consumes from specific shards it is assigned to. The ring matters for producers (where to enqueue) and for any component that looks up a tenant's queue, such as an admin tool or a depth metric per tenant.
How many virtual nodes should a ring use? 100–500 per shard gives even distribution for typical shard counts. More points cost memory and lookup time but improve balance.
Can Kafka use this to add partitions? Kafka's default partitioner is modulo-based, so adding partitions reshuffles keys. For keyed topics, over-provision partitions up front rather than adding them later — see Queue Partitioning Strategies.
Related
- Queue Partitioning Strategies — when and how to shard queues.
- Partitioning Queues by Tenant ID — tenant keys and pinning heavy tenants.
- Message Ordering Guarantees — why moving keys threatens ordering.
- Benchmarking Redis Broker Throughput — deciding when a new shard is needed.