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.

How many keys move when a shard is added Going from four shards to five, modulo hashing reassigns about 80 percent of tenants because almost every hash modulo 5 differs from the hash modulo 4. Consistent hashing reassigns about 20 percent, only the tenants that the new shard takes over, and every other tenant keeps its shard. Tenants that move, 4 shards to 5 hash % N ~80% move: ordering breaks, mass migration consistent hash ~20%: only the new shard's share The minimum possible is 1/(N+1): consistent hashing gets there.

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.

Adding a shard to the ring A circle of hash values holds virtual-node points for four shards, interleaved around the ring. Keys map to the next point clockwise. Adding redis-4 inserts its own points, each taking over only the arc just before it from whichever shard owned it. Keys elsewhere on the ring keep their shard. New points take small arcs from neighbours blue and navy points: shards 0-3 (virtual nodes) orange points: new shard redis-4 only keys just before an orange point move

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
Drain, then switch, one tenant at a time For each moving tenant, the migration pins it to its old shard, pauses its enqueues, waits until its queue on the old shard is empty with nothing in flight, switches the override to the new shard, and resumes enqueues. Jobs before the switch finish on the old shard; jobs after start on the new one, so the tenant's order is preserved. Moving tenant t-42 from redis-1 to redis-4 pin to redis-1 pause enqueue drain redis-1 override: redis-4 resume enqueue The pause lasts only as long as the tenant's backlog takes to drain, usually seconds. Every other tenant keeps running throughout.

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