Table of Contents

    consistent hashing

    DISTRIBUTED DATA & SHARDING

    Consistent Hashing

    Learn how consistent hashing distributes keys across a changing set of servers while minimizing remapping, how the hash ring and clockwise ownership work, why virtual nodes improve balance, and how replication, weighted placement, membership versioning, hot keys, rebalancing, cache warming, failure handling, and safe data migration affect production systems.

    Introduction

    Distributed caches, databases, storage systems, and request routers need a repeatable way to determine which server owns a given key.

    A simple solution is to hash the key and divide the result by the number of servers:

    \[ ServerIndex = Hash(Key) \bmod NumberOfServers \]

    This approach works while the server count remains unchanged.

    Three servers:
    
    server = hash(key) mod 3
    
    
    After adding another server:
    
    server = hash(key) mod 4

    Changing the divisor changes the result for many keys. Consequently, adding or removing one server can remap a large portion of the dataset.

    In a cache, widespread remapping can create many cache misses and send a sudden surge of requests to the database. In a sharded datastore, it can require extensive data movement before the new ownership arrangement becomes usable.

    Consistent hashing solves this problem by giving keys and servers positions in the same logical hash space. When membership changes, only the ownership ranges affected by that change need to move.

    Core idea: Consistent hashing minimizes key remapping when servers are added or removed. It does not eliminate rebalancing, but it limits the part of the keyspace whose ownership changes.

    Prerequisites

    # Prerequisite Why It Is Needed
    1 Hash functions Keys and server identifiers are mapped into a common hash space.
    2 Hash sharding Consistent hashing is a stable alternative to direct modulo-based placement.
    3 Database sharding The technique can assign data partitions or key ranges to database nodes.
    4 Distributed caching Cache nodes frequently join, leave, fail, and recover.
    5 Replication Each ownership range can require copies on additional nodes.
    6 Health checks and membership Routing changes when a node becomes available or unavailable.
    7 Observability Distribution, movement, hotspots, failures, and routing versions must be monitored.

    What Is Consistent Hashing?

    Consistent hashing is a key-placement technique designed to minimize reassignment when the set of available servers changes.

    Keys and servers are mapped into the same circular hash space, commonly called a hash ring.

    Logical hash ring:
    
                 0
                 |
           Node A|       Node B
              \  |  /
               \ | /
                \|/
    Maximum -----+----- Minimum
                /|\
               / | \
              /  |  \
           Node D       Node C

    To place a key:

    1. Hash the key into a ring position.
    2. Move clockwise from that position.
    3. Select the first node position encountered.
    4. Assign the key to the physical node represented by that position.
    Consistent-hashing Flow
    hash the key → locate ring position → move clockwise → select first eligible node → route operation

    The Problem with Modulo Hashing

    Direct modulo hashing selects a server using:

    \[ ServerIndex = Hash(Key) \bmod N \]

    Here, \(N\) is the number of servers.

    Before Adding a Server

    Number of servers:
    
    N = 3
    
    
    Key placement:
    
    hash(course:42) mod 3
    hash(user:1042) mod 3
    hash(session:981) mod 3

    After Adding a Server

    Number of servers:
    
    N = 4
    
    
    New placement:
    
    hash(course:42) mod 4
    hash(user:1042) mod 4
    hash(session:981) mod 4

    Because the divisor changed, many keys can now map to different server indexes.

    Consequences

    • Many cache entries become unreachable through their previous mapping
    • Database or storage records require data movement
    • Network traffic increases during redistribution
    • Backend databases can receive a surge of cache-miss traffic
    • Client instances can disagree if they update membership at different times

    Modulo Hashing vs Consistent Hashing

    Area Modulo Hashing Consistent Hashing
    Placement hash(key) mod nodeCount Hash key and select its clockwise owner on a ring
    Membership change Changes the divisor Changes selected ownership intervals
    Key movement Many keys can remap Movement is limited mainly to affected ranges
    Metadata Requires node count and stable indexing Requires ring membership and position metadata
    Unequal node capacity Difficult with simple modulo placement Can be modeled through weighted virtual-node ownership
    Operational complexity Simple until membership changes More placement metadata but safer incremental changes

    The Hash Ring

    A hash function produces values within a bounded numerical space.

    For example, a conceptual hash space could range from \(0\) to \(2^{32} - 1\). The highest value wraps around to zero, producing a logical ring.

    Hash space:
    
    0 -------------------------------- Maximum
    |                                      |
    +--------------------------------------+
                  wraps around

    Server identifiers are hashed onto this ring:

    hash("server-a") -> Position 10
    
    hash("server-b") -> Position 35
    
    hash("server-c") -> Position 68
    
    hash("server-d") -> Position 91

    Keys are hashed into the same space:

    hash("tenant:17") -> Position 30
    
    hash("course:42") -> Position 62
    
    hash("session:981") -> Position 96

    Clockwise Ownership

    A key belongs to the first eligible node encountered while moving clockwise from the key's position.

    Ring positions:
    
    Server A = 10
    Server B = 35
    Server C = 68
    Server D = 91
    
    
    Key positions:
    
    tenant:17 = 30
        -> Server B
    
    
    course:42 = 62
        -> Server C
    
    
    session:981 = 96
        -> Wrap around
        -> Server A

    The wrap-around rule makes the ring continuous.

    Ownership Intervals

    Each node owns an interval between its predecessor and its own position.

    Node Conceptual Ownership
    Server A at 10 Values after 91 through 10, including ring wrap-around
    Server B at 35 Values after 10 through 35
    Server C at 68 Values after 35 through 68
    Server D at 91 Values after 68 through 91

    Adding a Node

    Suppose a new server is added at ring position 50.

    Before:
    
    Server B = 35
    Server C = 68
    
    
    Server C owns:
    
    Values after 35 through 68
    
    
    After adding Server E at 50:
    
    Server E owns:
    
    Values after 35 through 50
    
    
    Server C continues to own:
    
    Values after 50 through 68

    Only keys in the newly assigned interval need to move from Server C to Server E.

    Add Server E
          |
          v
    Determine its predecessor
    and successor
          |
          v
    Identify newly owned interval
          |
          v
    Move affected keys
          |
          v
    Update membership version
          |
          v
    Enable normal routing

    Removing or Losing a Node

    When a node leaves the ring, its ownership interval transfers to its clockwise successor.

    Before failure:
    
    Server B owns:
    
    Values after 10 through 35
    
    
    Server B becomes unavailable
          |
          v
    Server C becomes responsible
    for the affected interval

    If the system uses replication, another replica can serve the affected partition while ownership is repaired. Without replication, the mapping can select another node but the actual data may remain unavailable.

    Availability rule: Consistent hashing determines placement. It does not create replicas or recover lost data automatically. Availability requires an explicit replication and repair design.

    Expected Data Movement

    In a well-balanced ring with \(N\) similarly sized nodes, adding one node should move approximately the new node's fair share of the keyspace.

    A simplified estimate is:

    \[ ApproximateMovedFraction = \frac{1}{N + 1} \]

    This is an expectation for reasonably balanced placement. Actual movement depends on ring positions, virtual nodes, key distribution, node weights, and workload skew.

    Why One Ring Position per Server Is Not Enough

    If each physical server receives only one random position, ownership ranges can differ substantially in size.

    Uneven ring:
    
    Server A owns a very small interval
    
    Server B owns a medium interval
    
    Server C owns a very large interval

    Server C can receive much more data and traffic than the other nodes even when keys are uniformly distributed across the hash space.

    A node failure can also transfer its complete interval to one successor, producing a sudden concentrated load increase.

    Virtual Nodes

    A virtual node is a logical position on the hash ring assigned to a physical server.

    Instead of receiving one position, each physical server receives several positions spread across the ring.

    Physical Server A:
    
    A1
    A2
    A3
    A4
    
    
    Physical Server B:
    
    B1
    B2
    B3
    B4
    
    
    Physical Server C:
    
    C1
    C2
    C3
    C4

    Each virtual position owns its own small interval. The physical server owns the union of the intervals assigned to its virtual nodes.

    Benefits of Virtual Nodes

    • Ownership is spread across many smaller ranges
    • Distribution is less dependent on a few random physical-node positions
    • Node failure transfers ranges to several successors
    • Rebalancing work can be distributed among several existing nodes
    • Different server capacities can receive different numbers of virtual positions

    Costs of Virtual Nodes

    • More ring metadata must be stored and distributed
    • Routing tables become larger
    • Membership changes can affect several small ranges
    • Monitoring must aggregate virtual ownership by physical node
    • Too many virtual nodes can increase management overhead

    Weighted Consistent Hashing

    Physical nodes can have different CPU, memory, storage, or network capacities.

    Server A capacity:
    
    1 unit
    
    
    Server B capacity:
    
    2 units
    
    
    Possible weighted placement:
    
    Server A receives X virtual nodes
    
    Server B receives approximately 2X virtual nodes

    Weighted placement gives higher-capacity nodes a larger expected share of the keyspace.

    Capacity is not always represented by one simple weight. A storage-heavy workload and a CPU-heavy workload can require different weighting models. Validate weights through representative load tests.

    Consistent Hashing Does Not Solve Hot Keys

    Consistent hashing can distribute many different keys across nodes. It cannot divide one indivisible popular key automatically.

    Normal keys:
    
    1,000 requests per minute each
    
    
    One popular key:
    
    500,000 requests per minute
    
    
    Consistent hashing:
    
    Popular key still has one
    primary ownership location.

    Possible hot-key controls include:

    • Replicating the popular read value
    • Using local or edge caching
    • Request coalescing
    • Splitting a logical key into subkeys where safe
    • Applying rate or concurrency limits
    • Using a dedicated placement policy

    Distribution rule: Even key distribution does not guarantee even traffic distribution. Measure request rate, stored bytes, CPU, I/O, and latency separately for every node and ownership range.

    Replication on the Ring

    A partition can be replicated to several distinct physical nodes encountered while moving clockwise.

    Key hashes to Position 42
          |
          v
    Primary owner:
    
    First eligible physical node clockwise
    
    
    Replica 1:
    
    Next distinct physical node clockwise
    
    
    Replica 2:
    
    Next distinct physical node clockwise

    Replica placement should consider physical failure domains. Placing every replica on servers within one machine, rack, zone, or region can leave the data vulnerable to one shared failure.

    Physical-node Deduplication

    When virtual nodes are used, several consecutive virtual positions can belong to the same physical server. Replica selection should skip duplicate physical owners where the durability policy requires distinct nodes.

    Consistent Hashing vs Replication

    Area Consistent Hashing Replication
    Primary question Which node owns this key? Which additional nodes store copies?
    Membership change Limits ownership reassignment Restores the required copy count
    Node failure Identifies a new routing owner Provides another copy from which service can continue
    Data durability Not provided by placement alone Depends on copy count and acknowledgment policy
    Repair Defines new ownership ranges Copies missing data to replacement owners

    Membership Management

    Every router must know the current set of nodes and virtual positions.

    Membership metadata can include:

    • Physical node identifier
    • Virtual-node positions
    • Node weight
    • Node health and lifecycle state
    • Failure-domain information
    • Membership version or epoch
    • Migration state

    If clients use different membership views, they can route the same key to different nodes.

    Client A ring version:
    
    42
    
    
    Client B ring version:
    
    41
    
    
    Same key:
    
    Client A -> Node D
    
    Client B -> Node C

    Ring Versions and Epochs

    A ring version identifies a specific membership and ownership configuration.

    {
      "ringVersion": 42,
      "hashAlgorithm": "approved-hash",
      "nodes": [
        {
          "nodeId": "cache-a",
          "state": "active",
          "weight": 1
        },
        {
          "nodeId": "cache-b",
          "state": "active",
          "weight": 2
        }
      ]
    }

    Requests, routing errors, and migration events can include the ring version without exposing sensitive infrastructure details to untrusted clients.

    Protecting Membership Metadata

    Membership configuration controls where data and requests are sent.

    Protect it with:

    • Authenticated configuration distribution
    • Least-privilege update permissions
    • Version validation
    • Private control-plane access
    • Audited membership changes
    • Rollback support
    • Backup and recovery

    Do not accept a public client's arbitrary node address, ring position, or membership version as authoritative routing information.

    Health Checks and Ring Membership

    A brief probe failure should not necessarily rewrite ownership immediately.

    One health probe fails
          |
          v
    If node is removed immediately:
    
    Ownership moves
          |
          v
    Next probe succeeds
          |
          v
    Ownership moves back
          |
          v
    Repeated rebalance and cache churn

    Use appropriate failure thresholds, lifecycle states, and stabilization policies.

    Node State Possible Meaning
    Joining The node is receiving or warming assigned ranges
    Active The node is eligible for ordinary routing
    Draining The node serves existing ownership while ranges move away
    Unavailable The node cannot serve requests
    Removed The node no longer belongs to the active ring

    Rebalancing

    Rebalancing moves affected keys or partitions to match a new ownership map.

    Create new ring version
          |
          v
    Identify changed ownership ranges
          |
          v
    Prepare destination nodes
          |
          v
    Copy or warm data
          |
          v
    Capture concurrent changes
          |
          v
    Validate destination
          |
          v
    Shift request routing
          |
          v
    Retire old ownership safely

    Consistent hashing limits how much ownership changes. It does not define the complete data-migration protocol.

    Rebalancing Strategies

    Strategy Direction Main Risk
    Cold reassignment Change ownership and allow the new cache node to populate naturally Cache misses can increase backend load
    Pre-warming Copy or load likely-needed data before enabling traffic Warm data can become stale during preparation
    Streaming migration Transfer data from the old owner to the new owner Migration consumes network and storage capacity
    Dual-read transition Read the new owner and fall back to the old owner temporarily Temporary complexity and additional read traffic
    Dual-write transition Write to old and new owners during migration Partial success can create divergent copies

    Cache Rebalancing

    A distributed cache can tolerate losing cached entries because the authoritative value remains in another datastore. However, many simultaneous misses can overload that datastore.

    New cache node joins
          |
          v
    Some keys move to new owner
          |
          v
    Moved keys are initially absent
          |
          v
    Requests miss cache
          |
          v
    Database load increases

    Protect the origin using:

    • Gradual traffic shifting
    • Cache pre-warming
    • Request coalescing
    • Bounded concurrency
    • Rate limiting
    • Approved stale-value serving
    • Autoscaling headroom

    Database Rebalancing

    Database placement requires stronger migration controls because the moved data is authoritative.

    Design for:

    • Historical copy
    • Concurrent writes during migration
    • Version validation
    • Checksums or reconciliation
    • Atomic or coordinated routing cutover
    • Rollback
    • Replica restoration
    • Retention of the old copy until validation completes

    Migration rule: Changing the ring does not move authoritative data safely by itself. Use a migration protocol that handles concurrent writes, validation, cutover, retries, and rollback.

    Routing Implementation

    A router can store virtual-node positions in a sorted collection.

    For each key:

    1. Calculate the key hash.
    2. Search for the first ring position greater than or equal to that hash.
    3. If no position exists, wrap to the first position.
    4. Resolve the virtual position to its physical node.

    Binary search can locate the responsible virtual position efficiently in a sorted ring.

    PHP Consistent-hash Router

    <?php
    
    declare(strict_types=1);
    
    final class ConsistentHashRing
    {
        /**
         * @var array<int, string>
         */
        private array $ring = [];
    
        /**
         * @var list<int>
         */
        private array $positions = [];
    
        public function addNode(
            string $nodeId,
            int $virtualNodeCount
        ): void {
            if ($virtualNodeCount < 1) {
                throw new InvalidArgumentException(
                    'Virtual-node count must be positive.'
                );
            }
    
            for (
                $index = 0;
                $index < $virtualNodeCount;
                $index++
            ) {
                $virtualNodeId =
                    $nodeId . '#' . $index;
    
                $position =
                    $this->hash(
                        $virtualNodeId
                    );
    
                $this->ring[$position] =
                    $nodeId;
            }
    
            $this->rebuildPositions();
        }
    
        public function removeNode(
            string $nodeId
        ): void {
            foreach (
                $this->ring
                as $position => $currentNodeId
            ) {
                if ($currentNodeId === $nodeId) {
                    unset(
                        $this->ring[$position]
                    );
                }
            }
    
            $this->rebuildPositions();
        }
    
        public function findNode(
            string $key
        ): string {
            if ($this->positions === []) {
                throw new RuntimeException(
                    'The hash ring has no nodes.'
                );
            }
    
            $keyPosition =
                $this->hash(
                    $key
                );
    
            foreach (
                $this->positions
                as $position
            ) {
                if ($position >= $keyPosition) {
                    return $this->ring[
                        $position
                    ];
                }
            }
    
            $firstPosition =
                $this->positions[0];
    
            return $this->ring[
                $firstPosition
            ];
        }
    
        private function rebuildPositions(): void
        {
            $this->positions =
                array_keys(
                    $this->ring
                );
    
            sort(
                $this->positions,
                SORT_NUMERIC
            );
        }
    
        private function hash(
            string $value
        ): int {
            return (int)sprintf(
                '%u',
                crc32($value)
            );
        }
    }

    This is an educational example. A production implementation must use the platform's approved hash algorithm, collision handling, membership distribution, node health policy, replication strategy, ring versioning, migration process, and operational safeguards.

    Hash Collisions

    Two virtual-node identifiers can theoretically produce the same hash position.

    A production implementation should define:

    • How collisions are detected
    • How positions are made unique
    • Whether another salt or sequence is used
    • How deterministic membership is preserved across routers

    Do not silently overwrite one node's placement entry when a ring-position collision occurs.

    Choosing the Routing Key

    Consistent hashing can distribute only the key provided to it.

    Possible routing keys include:

    • Tenant ID
    • User ID
    • Session ID
    • Cache key
    • Object identifier
    • Partition identifier

    The routing key should have enough distinct values and should align with the required data locality.

    Routing-key Examples

    Key Benefit Risk
    Tenant ID Tenant data can remain colocated One very large tenant can overload one ownership location
    User ID User-level distribution can be broad Tenant-wide operations can fan out
    Session ID Sessions can distribute across cache nodes Related account sessions can be spread widely
    Object ID Point retrieval is directly routable Related objects might not be colocated
    Constant or low-cardinality key Simple routing Creates severe concentration and poor distribution

    Data Locality

    Records that must frequently be accessed or updated together should use a routing design that keeps the operation local where practical.

    Hash by tenant ID:
    
    Tenant profile
    Tenant users
    Tenant courses
    Tenant enrollments
    Tenant progress
    
    All reach the same tenant owner.

    Hashing each individual record identifier independently can distribute data more broadly but can turn a tenant-scoped transaction into a multi-node operation.

    Scatter-Gather Queries

    A query that does not contain the routing key can require requests to several or all owners.

    Global query
          |
          +-- Query Node A
          +-- Query Node B
          +-- Query Node C
          +-- Query Node D
          |
          v
    Merge partial results

    Consistent hashing is best suited to point operations that contain the placement key. Frequent global queries can require:

    • A separate search index
    • A directory service
    • An analytical projection
    • A global lookup table
    • A controlled scatter-gather layer

    Consistent Hashing vs Directory Sharding

    Area Consistent Hashing Directory Sharding
    Placement rule Derived from the ring and hash Stored explicitly for each key or tenant
    Lookup Calculated locally from membership metadata Requires directory lookup or cache
    Selective movement Normally moves ownership ranges Can move individual tenants or keys
    Large tenant isolation Requires special key splitting or placement support Tenant can be mapped to a dedicated shard
    Membership change Limits affected ring intervals Mappings are updated explicitly
    Critical dependency Membership metadata and consistent router view Highly available directory service

    Consistent Hashing vs Rendezvous Hashing

    Rendezvous hashing, also called highest-random-weight hashing, is another stable placement method.

    For each key:
    
    1. Calculate a score for
       every eligible node.
    
    2. Select the node with
       the highest score.
    Area Consistent Hashing Rendezvous Hashing
    Placement structure Logical ring with ordered positions Per-key score across eligible nodes
    Lookup Locate the clockwise owner Calculate and compare node scores
    Metadata Ring positions and virtual nodes Eligible node list and scoring rule
    Membership stability Moves affected intervals Reassigns keys associated with the changed node

    Product support, node count, lookup cost, weighting, replication, and operational tooling should guide the selection.

    Common Use Cases

    Use Case Placement Key Primary Benefit
    Distributed cache Cache key Limits cache-key remapping during node changes
    Sharded key-value store Record or partition key Distributes ownership across storage nodes
    Session store Session identifier Distributes sessions among cache or storage servers
    Object storage Object identifier Selects storage ownership consistently
    Message partitioning Message or entity key Keeps related messages on a stable owner
    Request routing Tenant, account, or resource key Provides stable backend affinity with reduced reassignment

    Learning-platform Example

    Session request
          |
          v
    Extract protected session ID
          |
          v
    Hash session ID
          |
          v
    Find clockwise cache owner
          |
          v
    Read session data

    Example Workloads

    Data Possible Routing Key Main Consideration
    Authenticated sessions Opaque session identifier Replicate sessions or retain an authoritative fallback
    Course-content cache Course ID Popular courses can become hot keys
    Learner-progress cache Tenant and learner identifier Preserve tenant isolation and invalidation correctness
    Video-processing jobs Asset ID Keep related processing on a stable partition where ordering matters
    Tenant database placement Tenant ID Large tenants can require dedicated placement rather than ordinary hashing

    Conceptual Ring Configuration

    consistentHashing:
      hashAlgorithm: approved-stable-hash
      ringVersion: 42
    
      virtualNodes:
        enabled: true
        allocation: capacity-weighted
    
      nodes:
        - nodeId: cache-a
          weight: 1
          state: active
    
        - nodeId: cache-b
          weight: 2
          state: active
    
        - nodeId: cache-c
          weight: 1
          state: joining
    
      ownership:
        select: clockwise-successor
    
      replication:
        enabled: true
        distinctPhysicalNodesRequired: true
    
      migration:
        preWarm: true
        validateBeforeActivation: true
        gradualTrafficShift: true
        rollbackRequired: true
    
      observability:
        ownershipDistribution: enabled
        hotKeyDetection: enabled
        ringVersionMismatch: enabled
        migrationProgress: enabled

    This is conceptual configuration. Use the actual hashing, placement, replication, migration, and membership facilities of the selected platform.

    Observability

    Useful consistent-hashing metrics include:

    • Physical-node count
    • Virtual-node count
    • Ring or membership version
    • Owned hash ranges by node
    • Stored bytes by node
    • Key count by node
    • Request rate by node
    • CPU, memory, network, and I/O by node
    • Hot-key frequency
    • Cache hit rate by node
    • Keys or partitions moved during rebalancing
    • Migration throughput and backlog
    • Routing-version mismatches
    • Wrong-owner responses
    • Replication health by ownership range
    • Node join and removal duration

    Structured Routing Event

    {
      "routingStrategy": "consistent-hashing",
      "ringVersion": 42,
      "virtualNodePosition": 185720,
      "physicalNode": "cache-b",
      "operation": "read",
      "result": "routed"
    }

    Avoid recording complete session identifiers, credentials, private record values, or confidential infrastructure details in routine routing logs.

    Alert Conditions

    Alert when:

    • One physical node owns a disproportionate keyspace share
    • Request traffic is substantially skewed despite balanced key ownership
    • One key dominates node traffic
    • Routers use different ring versions
    • A joining node cannot complete warming or migration
    • A draining node continues receiving new ownership traffic
    • Cache misses increase sharply after a membership change
    • Origin database traffic increases during rebalancing
    • Replication falls below the required copy count
    • Wrong-owner responses increase
    • A failed node's ranges cannot be reassigned or restored
    • Migration backlog threatens the next topology change

    Troubleshooting Workflow

    1. Identify the key used for routing.
    2. Calculate or inspect the key's hash position.
    3. Identify the expected clockwise virtual-node owner.
    4. Resolve the virtual node to its physical server.
    5. Compare ring versions across clients and routers.
    6. Check node lifecycle and health states.
    7. Check whether a membership change is in progress.
    8. Check ownership distribution and virtual-node weights.
    9. Check for hot keys or hot tenants.
    10. Check whether replication copies remain available.
    11. Check migration, warming, and repair progress.
    12. Check backend load caused by cache misses.
    13. Check recent changes to the hash algorithm or node identifiers.
    14. Roll back or repair membership through the approved procedure.

    Common Consistent-hashing Mistakes

    1

    Using Direct Modulo Hashing for a Dynamic Cluster

    Changing the node count remaps many keys and can create substantial cache churn or data movement.

    2

    Assigning Only One Position per Physical Node

    Random ring spacing can produce highly uneven ownership ranges.

    3

    Assuming Virtual Nodes Solve Hot Keys

    A very popular indivisible key remains assigned to one primary owner.

    4

    Changing the Hash Algorithm without Migration

    Existing keys can suddenly resolve to different ring positions and owners.

    5

    Allowing Routers to Use Different Membership Versions

    The same key can be sent to different physical nodes.

    6

    Removing a Node after One Failed Probe

    Temporary probe variation can trigger repeated ownership movement and cache churn.

    7

    Routing Traffic before the New Node Is Ready

    Requests reach an empty cache or incomplete database partition before warming and migration finish.

    8

    Changing the Ring before Authoritative Data Moves

    The new owner receives requests before the required records are available.

    9

    Using Unsafe Dual Writes during Migration

    One owner can accept the write while the other fails, creating divergent data.

    10

    Ignoring Physical Failure Domains

    Primary and replica ranges can be placed on infrastructure sharing one failure boundary.

    11

    Logging Sensitive Routing Keys

    Session identifiers, user identifiers, or confidential tenant information can be exposed through routing diagnostics.

    12

    Testing Distribution but Not Rebalancing

    Membership changes, backend surges, migration failures, and rollback remain unverified.

    Recommended Test Cases

    Test Expected Evidence
    Deterministic routing The same key and ring version always select the same owner
    Ring wrap-around A key beyond the final ring position routes to the first eligible node
    Add one node Only the ownership ranges assigned to the new node move
    Remove one node Only the removed node's ownership ranges are reassigned
    Virtual-node distribution Ownership is distributed within the approved balance objective
    Weighted placement Higher-capacity nodes receive the intended larger ownership share
    Hot key One popular key is detected and protected separately
    Node failure Replicas or fallback owners serve affected ranges according to policy
    Temporary health failure The ring does not flap because of one brief probe failure
    Stale ring version The router refreshes membership and performs only a bounded retry
    Cache-node addition Origin load remains within its safe capacity while keys warm
    Database migration Historical and concurrent writes exist correctly on the new owner
    Migration rollback Routing returns safely to the previous ownership map
    Replication placement Copies are stored on distinct eligible physical failure domains
    Hash-algorithm compatibility Every router calculates identical positions for known test keys

    Consistent-hashing Best Practices

    Recommended Practices

    • Use consistent hashing when cluster membership changes are expected.
    • Use a stable, deterministic, approved hash function.
    • Hash keys and node positions into the same space.
    • Route each key to the first eligible clockwise owner.
    • Use virtual nodes to improve distribution.
    • Weight virtual-node ownership when physical capacity differs.
    • Measure key, storage, and traffic distribution separately.
    • Detect and protect hot keys explicitly.
    • Version all membership and ring configurations.
    • Ensure routers converge on the same membership view.
    • Use stabilized health decisions before changing membership.
    • Warm or migrate data before enabling full traffic.
    • Protect origin systems from cache-miss storms.
    • Replicate ownership ranges across distinct physical nodes.
    • Consider rack, zone, and regional failure domains.
    • Keep frequent operations local by choosing an appropriate routing key.
    • Use a separate lookup or projection for non-key queries.
    • Use a safe protocol for authoritative data movement.
    • Support rollback for membership and ownership changes.
    • Test joins, removals, failures, stale rings, hotspots, and migrations.

    Practice Exercise

    Design consistent hashing for your online learning platform's distributed session and course-content cache.

    Requirements

    1. Select a stable hash function.
    2. Define the logical hash space.
    3. Place at least three physical cache nodes on the ring.
    4. Assign multiple virtual nodes to each physical server.
    5. Route session identifiers to their clockwise owners.
    6. Route course-content keys separately.
    7. Add replication across distinct physical nodes.
    8. Add a fourth cache server.
    9. Measure which ownership ranges move.
    10. Warm moved cache entries without overloading the database.
    11. Remove one server and verify replica routing.
    12. Introduce one hot course key.
    13. Add request coalescing and hot-key protection.
    14. Create and distribute a versioned ring configuration.
    15. Test a client using a stale ring version.
    16. Test rollback after a failed membership change.
    17. Monitor key count, traffic, memory, hit rate, and movement by node.

    Design Template

    Decision Selected Direction Primary Risk Controlled
    Routing key Protected session ID or stable resource ID Deterministic point routing
    Placement Clockwise owner on a versioned hash ring Limits remapping during membership changes
    Distribution Several virtual nodes per physical server Reduces uneven ownership ranges
    Capacity differences Capacity-weighted virtual-node allocation Prevents weaker nodes from receiving equal ownership
    Availability Replica placement on distinct physical nodes Allows service during one node failure
    Node addition Warm, validate, then activate moved ranges Prevents origin surge and premature routing
    Hot key Replication, coalescing, and bounded concurrency Protects one owner from concentrated demand
    Membership agreement Authenticated versioned ring distribution Prevents inconsistent routing among clients

    Frequently Asked Questions

    1

    What is consistent hashing?

    Consistent hashing is a placement technique that limits key reassignment when servers are added or removed.

    2

    What problem does consistent hashing solve?

    It avoids the widespread remapping caused by changing the divisor in direct hash(key) mod nodeCount placement.

    3

    What is a hash ring?

    A hash ring is a circular representation of the hash space on which keys and node positions are placed.

    4

    How is a key assigned?

    The key is hashed to a ring position and assigned to the first eligible node encountered clockwise.

    5

    What happens when a node is added?

    The new node takes ownership of selected intervals previously owned by its clockwise successor.

    6

    What happens when a node is removed?

    The removed node's intervals move to their next eligible clockwise owners.

    7

    What is a virtual node?

    A virtual node is one logical ring position assigned to a physical server. A server normally owns several such positions.

    8

    Why are virtual nodes useful?

    They divide ownership into smaller ranges, improve balance, distribute failure load, and support capacity-weighted placement.

    9

    Does consistent hashing solve hot keys?

    No. It distributes many keys, but one highly popular key remains concentrated on its primary owner unless additional controls are used.

    10

    Does consistent hashing replicate data?

    No. It determines ownership. Replication must be designed separately using additional eligible physical nodes.

    11

    Does consistent hashing eliminate data movement?

    No. It limits ownership changes, but affected cache entries or database partitions still need warming, copying, repair, or migration.

    12

    Where is consistent hashing used?

    It is commonly applied to distributed caches, key-value stores, sharded storage, session stores, object placement, message partitioning, and stable request routing.

    Key Takeaway

    Consistent hashing maps keys and servers into one logical hash ring and assigns each key to its first eligible clockwise owner. Unlike direct modulo hashing, adding or removing a node changes only selected ownership intervals rather than remapping much of the keyspace. Virtual nodes improve distribution by giving every physical server several smaller intervals, and weighted placement can assign more ownership to higher-capacity servers. Consistent hashing does not solve hot keys, create replicas, move authoritative data safely, or keep routers synchronized automatically. Production designs still require versioned membership, stable hash functions, health-state stabilization, replication across failure domains, safe migration, cache warming, origin protection, hotspot detection, and rollback. Test node joins, removals, failures, stale ring versions, rebalancing, and backend-loading effects before depending on the ring in production.