From 0330a8153fafb9d70cdca2dbad7a3a7b19542d79 Mon Sep 17 00:00:00 2001 From: PJ Fanning Date: Fri, 9 Oct 2026 16:13:33 +0100 Subject: [PATCH 1/2] fix: keep colliding virtual nodes in ConsistentHash Motivation: `ConsistentHash` stores ring positions in a `SortedMap[Int, T]`. When virtual nodes of two different nodes hash to the same position, the later entry silently overwrites the earlier one. The ring then depends on the order in which nodes were given or added, so two instances built from the same nodes can route the same key to different nodes. Removing a node also deletes any position it collides on, even if that position is owned by another node (or the removed node was never in the ring), so the remaining node loses virtual nodes. Modification: - Keep all nodes that claim a ring position: the owner in the ring, the others in a `collisions` map, which is normally empty. - On a collision the node with the lowest `toString` owns the position, independent of insertion order. - `:-` only releases the removed node's claims and promotes the next claimant if the owner is removed. - Document the collision rule in the class Scaladoc. - Add `ConsistentHashSpec` using two node names whose virtual nodes are known to collide. Result: The ring, and so `nodeFor`, is the same regardless of node order, and `ch :- node` routes like a ring built without that node. Routing only changes compared to before for keys that land on a colliding position. Tests: - `actor-tests/Test/testOnly org.apache.pekko.routing.ConsistentHashSpec` without the fix: 5 of 6 failed - `actor-tests/Test/testOnly org.apache.pekko.routing.ConsistentHashSpec org.apache.pekko.routing.ConsistentHashingRouterSpec` with the fix: 9 passed on 2.13.18 and 3.3.8 - `actor/mimaReportBinaryIssues`: passed on 2.13.18 and 3.3.8 - scalafmt run on changed files References: Fixes #3279 --- .../pekko/routing/ConsistentHashSpec.scala | 76 +++++++++++++++++ .../apache/pekko/routing/ConsistentHash.scala | 84 +++++++++++++++---- 2 files changed, 142 insertions(+), 18 deletions(-) create mode 100644 actor-tests/src/test/scala/org/apache/pekko/routing/ConsistentHashSpec.scala diff --git a/actor-tests/src/test/scala/org/apache/pekko/routing/ConsistentHashSpec.scala b/actor-tests/src/test/scala/org/apache/pekko/routing/ConsistentHashSpec.scala new file mode 100644 index 00000000000..b6b4af750d0 --- /dev/null +++ b/actor-tests/src/test/scala/org/apache/pekko/routing/ConsistentHashSpec.scala @@ -0,0 +1,76 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ + +package org.apache.pekko.routing + +import org.scalatest.matchers.should.Matchers +import org.scalatest.wordspec.AnyWordSpec + +class ConsistentHashSpec extends AnyWordSpec with Matchers { + + // virtual nodes of these two nodes hash to the same ring positions with a factor of 10 + private val nodeA = "node-4230" + private val nodeB = "node-14323" + private val nodeC = "node-1" + private val factor = 10 + private val keys = (0 until 10000).map(i => s"key-$i") + + private def routing(ch: ConsistentHash[String]): Map[String, String] = + keys.iterator.map(key => key -> ch.nodeFor(key)).toMap + + "ConsistentHash" must { + + "route keys independently of the order the nodes were given" in { + routing(ConsistentHash(List(nodeA, nodeB, nodeC), factor)) should + ===(routing(ConsistentHash(List(nodeB, nodeA, nodeC), factor))) + routing(ConsistentHash(List(nodeA, nodeB, nodeC), factor)) should + ===(routing(ConsistentHash(List(nodeC, nodeB, nodeA), factor))) + } + + "route keys independently of the order the nodes were added" in { + val expected = routing(ConsistentHash(List(nodeA, nodeB, nodeC), factor)) + routing(ConsistentHash(List(nodeC), factor) :+ nodeA :+ nodeB) should ===(expected) + routing(ConsistentHash(List(nodeC), factor) :+ nodeB :+ nodeA) should ===(expected) + } + + "restore the colliding virtual nodes of the remaining node when a node is removed" in { + routing(ConsistentHash(List(nodeA, nodeB, nodeC), factor) :- nodeA) should + ===(routing(ConsistentHash(List(nodeB, nodeC), factor))) + routing(ConsistentHash(List(nodeB, nodeA, nodeC), factor) :- nodeA) should + ===(routing(ConsistentHash(List(nodeB, nodeC), factor))) + routing(ConsistentHash(List(nodeA, nodeB, nodeC), factor) :- nodeB) should + ===(routing(ConsistentHash(List(nodeA, nodeC), factor))) + routing(ConsistentHash(List(nodeB, nodeA, nodeC), factor) :- nodeB) should + ===(routing(ConsistentHash(List(nodeA, nodeC), factor))) + } + + "not remove virtual nodes of other nodes when removing a node that is not in the ring" in { + routing(ConsistentHash(List(nodeB, nodeC), factor) :- nodeA) should + ===(routing(ConsistentHash(List(nodeB, nodeC), factor))) + } + + "be empty after all nodes are removed" in { + (ConsistentHash(List(nodeA, nodeB), factor) :- nodeA :- nodeB).isEmpty should ===(true) + } + + "not change the ring when a node is added again" in { + val ch = ConsistentHash(List(nodeA, nodeB, nodeC), factor) + routing(ch :+ nodeA) should ===(routing(ch)) + routing(ch :+ nodeB) should ===(routing(ch)) + } + } +} diff --git a/actor/src/main/scala/org/apache/pekko/routing/ConsistentHash.scala b/actor/src/main/scala/org/apache/pekko/routing/ConsistentHash.scala index 343970fa552..4ab75408022 100644 --- a/actor/src/main/scala/org/apache/pekko/routing/ConsistentHash.scala +++ b/actor/src/main/scala/org/apache/pekko/routing/ConsistentHash.scala @@ -26,8 +26,17 @@ import scala.reflect.ClassTag * * Note that toString of the ring nodes are used for the node * hash, i.e. make sure it is different for different nodes. + * + * If virtual nodes of different nodes hash to the same ring position, the node + * with the lowest toString owns that position, so the ring does not depend on the + * order in which nodes were added. */ -class ConsistentHash[T: ClassTag] private (nodes: immutable.SortedMap[Int, T], val virtualNodesFactor: Int) { +class ConsistentHash[T: ClassTag] private ( + nodes: immutable.SortedMap[Int, T], + // other nodes whose virtual nodes hash to an owned ring position, kept so that + // they can take over the position if its owner is removed + collisions: immutable.Map[Int, List[T]], + val virtualNodesFactor: Int) { import ConsistentHash._ @@ -47,11 +56,8 @@ class ConsistentHash[T: ClassTag] private (nodes: immutable.SortedMap[Int, T], v * operation returns a new instance. */ def :+(node: T): ConsistentHash[T] = { - val nodeHash = hashFor(node.toString) - new ConsistentHash(nodes ++ - ((1 to virtualNodesFactor).map { r => - concatenateNodeHash(nodeHash, r) -> node - }), virtualNodesFactor) + val (newNodes, newCollisions) = claim(nodes, collisions, node, virtualNodesFactor) + new ConsistentHash(newNodes, newCollisions, virtualNodesFactor) } /** @@ -68,10 +74,13 @@ class ConsistentHash[T: ClassTag] private (nodes: immutable.SortedMap[Int, T], v */ def :-(node: T): ConsistentHash[T] = { val nodeHash = hashFor(node.toString) - new ConsistentHash(nodes -- - ((1 to virtualNodesFactor).map { r => - concatenateNodeHash(nodeHash, r) - }), virtualNodesFactor) + val (newNodes, newCollisions) = (1 to virtualNodesFactor).foldLeft((nodes, collisions)) { + case ((ns, cs), r) => + val hash = concatenateNodeHash(nodeHash, r) + val others = claimants(ns, cs, hash).filterNot(sameNode(_, node)) + setClaimants(ns, cs, hash, others) + } + new ConsistentHash(newNodes, newCollisions, virtualNodesFactor) } /** @@ -123,14 +132,11 @@ class ConsistentHash[T: ClassTag] private (nodes: immutable.SortedMap[Int, T], v object ConsistentHash { def apply[T: ClassTag](nodes: Iterable[T], virtualNodesFactor: Int): ConsistentHash[T] = { - new ConsistentHash( - immutable.SortedMap.empty[Int, T] ++ - (for { - node <- nodes - nodeHash = hashFor(node.toString) - vnode <- 1 to virtualNodesFactor - } yield concatenateNodeHash(nodeHash, vnode) -> node), - virtualNodesFactor) + val (ring, collisions) = + nodes.foldLeft((immutable.SortedMap.empty[Int, T], immutable.Map.empty[Int, List[T]])) { + case ((ns, cs), node) => claim(ns, cs, node, virtualNodesFactor) + } + new ConsistentHash(ring, collisions, virtualNodesFactor) } /** @@ -142,6 +148,48 @@ object ConsistentHash { apply(nodes.asScala, virtualNodesFactor) } + // nodes are identified by their toString, see the class documentation + private def sameNode[T](a: T, b: T): Boolean = a.toString == b.toString + + // all nodes with a virtual node at the given ring position, the owner first + private def claimants[T]( + nodes: immutable.SortedMap[Int, T], + collisions: immutable.Map[Int, List[T]], + hash: Int): List[T] = + nodes.get(hash) match { + case Some(owner) => owner :: collisions.getOrElse(hash, Nil) + case None => Nil + } + + private def setClaimants[T]( + nodes: immutable.SortedMap[Int, T], + collisions: immutable.Map[Int, List[T]], + hash: Int, + claimants: List[T]): (immutable.SortedMap[Int, T], immutable.Map[Int, List[T]]) = + claimants match { + case Nil => (nodes - hash, collisions - hash) + case single :: Nil => (nodes.updated(hash, single), collisions - hash) + case _ => + // lowest toString owns the position, independent of the order the nodes were added + val sorted = claimants.sortBy(_.toString) + (nodes.updated(hash, sorted.head), collisions.updated(hash, sorted.tail)) + } + + // adds the virtual nodes of `node`, replacing any previous virtual nodes of the same node + private def claim[T]( + nodes: immutable.SortedMap[Int, T], + collisions: immutable.Map[Int, List[T]], + node: T, + virtualNodesFactor: Int): (immutable.SortedMap[Int, T], immutable.Map[Int, List[T]]) = { + val nodeHash = hashFor(node.toString) + (1 to virtualNodesFactor).foldLeft((nodes, collisions)) { + case ((ns, cs), r) => + val hash = concatenateNodeHash(nodeHash, r) + val others = claimants(ns, cs, hash).filterNot(sameNode(_, node)) + setClaimants(ns, cs, hash, node :: others) + } + } + private def concatenateNodeHash(nodeHash: Int, vnode: Int): Int = { import MurmurHash._ var h = startHash(nodeHash) From 7f93376ae69a2bfaaeb33efb3f1c12c0dd22d675 Mon Sep 17 00:00:00 2001 From: PJ Fanning Date: Fri, 9 Oct 2026 16:53:16 +0100 Subject: [PATCH 2/2] perf: build the ConsistentHash ring in one sorted batch in apply Motivation: Resolving collisions by claiming one virtual node at a time made `ConsistentHash.apply` insert into the persistent ring map point by point. Against the previous bulk `SortedMap.empty ++ pairs` build that is 2-3x the allocation and several times the time (JMH, 10 virtual nodes per node: 100 nodes 230 KB -> 576 KB, 1000 nodes 2.4 MB -> 7.0 MB). `apply` is used to rebuild the ring on membership changes by the classic consistent hashing router, `ConsistentHashingShardAllocationStrategy` and `ClusterClient`. Modification: - Encode each virtual node as `(ring position << 32 | index)` in a primitive `Long` array, sort it once, and add positions with a single claimant through a `TreeMap` builder. - Resolve only runs of equal positions (collisions) with the same owner rule as `:+`, via a shared `ownerOf` helper. Result: `apply` produces the same ring as before this commit, and allocates less than the original (pre-fix) bulk build at similar speed. JMH, 10 virtual nodes per node, original vs this commit: - 10 nodes: 24.4 KB -> 18.5 KB, 12 us -> 22 us (+-20 us, noisy) - 100 nodes: 230 KB -> 201 KB, 219 us -> 235 us - 1000 nodes: 2.39 MB -> 1.87 MB, 3378 us -> 2899 us Tests: - `actor-tests/Test/testOnly org.apache.pekko.routing.ConsistentHashSpec org.apache.pekko.routing.ConsistentHashingRouterSpec`: 10 passed on 2.13.18 and 3.3.8, including a new check that `apply` matches adding the nodes one by one with `:+` - `actor/mimaReportBinaryIssues`: passed on 2.13.18 and 3.3.8 - scalafmt run on changed files References: Refs #3279 --- .../pekko/routing/ConsistentHashSpec.scala | 7 +++ .../apache/pekko/routing/ConsistentHash.scala | 63 ++++++++++++++++--- 2 files changed, 62 insertions(+), 8 deletions(-) diff --git a/actor-tests/src/test/scala/org/apache/pekko/routing/ConsistentHashSpec.scala b/actor-tests/src/test/scala/org/apache/pekko/routing/ConsistentHashSpec.scala index b6b4af750d0..cdd4b688677 100644 --- a/actor-tests/src/test/scala/org/apache/pekko/routing/ConsistentHashSpec.scala +++ b/actor-tests/src/test/scala/org/apache/pekko/routing/ConsistentHashSpec.scala @@ -63,6 +63,13 @@ class ConsistentHashSpec extends AnyWordSpec with Matchers { ===(routing(ConsistentHash(List(nodeB, nodeC), factor))) } + "build the same ring with apply as by adding the nodes one by one" in { + val many = (0 until 200).map(n => s"node-host-$n") ++ List(nodeA, nodeB, nodeC, nodeA) + val incremental = many.foldLeft(ConsistentHash(Nil: Seq[String], factor))(_ :+ _) + routing(ConsistentHash(many, factor)) should ===(routing(incremental)) + routing(ConsistentHash(many.reverse, factor)) should ===(routing(incremental)) + } + "be empty after all nodes are removed" in { (ConsistentHash(List(nodeA, nodeB), factor) :- nodeA :- nodeB).isEmpty should ===(true) } diff --git a/actor/src/main/scala/org/apache/pekko/routing/ConsistentHash.scala b/actor/src/main/scala/org/apache/pekko/routing/ConsistentHash.scala index 4ab75408022..355a1e673d9 100644 --- a/actor/src/main/scala/org/apache/pekko/routing/ConsistentHash.scala +++ b/actor/src/main/scala/org/apache/pekko/routing/ConsistentHash.scala @@ -132,11 +132,52 @@ class ConsistentHash[T: ClassTag] private ( object ConsistentHash { def apply[T: ClassTag](nodes: Iterable[T], virtualNodesFactor: Int): ConsistentHash[T] = { - val (ring, collisions) = - nodes.foldLeft((immutable.SortedMap.empty[Int, T], immutable.Map.empty[Int, List[T]])) { - case ((ns, cs), node) => claim(ns, cs, node, virtualNodesFactor) + if (virtualNodesFactor < 1) throw new IllegalArgumentException("virtualNodesFactor must be >= 1") + val nodeArray = nodes.toArray + val total = nodeArray.length * virtualNodesFactor + // each virtual node is encoded as (ring position << 32 | index), so that a primitive sort orders + // them by ring position, and by insertion order for the same position + val points = new Array[Long](total) + var n = 0 + while (n < nodeArray.length) { + val nodeHash = hashFor(nodeArray(n).toString) + var r = 1 + while (r <= virtualNodesFactor) { + val i = n * virtualNodesFactor + r - 1 + points(i) = (concatenateNodeHash(nodeHash, r).toLong << 32) | i + r += 1 } - new ConsistentHash(ring, collisions, virtualNodesFactor) + n += 1 + } + Arrays.sort(points) + + def hashAt(i: Int): Int = (points(i) >> 32).toInt + def nodeAt(i: Int): T = nodeArray((points(i) & 0xFFFFFFFFL).toInt / virtualNodesFactor) + + val ring = immutable.TreeMap.newBuilder[Int, T] + var collisions = immutable.Map.empty[Int, List[T]] + var i = 0 + while (i < total) { + val hash = hashAt(i) + var end = i + 1 + while (end < total && hashAt(end) == hash) end += 1 + if (end == i + 1) ring += hash -> nodeAt(i) + else { + // rare: several virtual nodes at the same ring position, resolve them like `:+` does + var claimants = List.empty[T] + var j = i + while (j < end) { + val node = nodeAt(j) + claimants = node :: claimants.filterNot(sameNode(_, node)) + j += 1 + } + val (owner, others) = ownerOf(claimants) + ring += hash -> owner + if (others.nonEmpty) collisions = collisions.updated(hash, others) + } + i = end + } + new ConsistentHash(ring.result(), collisions, virtualNodesFactor) } /** @@ -166,13 +207,19 @@ object ConsistentHash { collisions: immutable.Map[Int, List[T]], hash: Int, claimants: List[T]): (immutable.SortedMap[Int, T], immutable.Map[Int, List[T]]) = + if (claimants.isEmpty) (nodes - hash, collisions - hash) + else { + val (owner, others) = ownerOf(claimants) + (nodes.updated(hash, owner), if (others.isEmpty) collisions - hash else collisions.updated(hash, others)) + } + + // the lowest toString owns the position, independent of the order the nodes were added + private def ownerOf[T](claimants: List[T]): (T, List[T]) = claimants match { - case Nil => (nodes - hash, collisions - hash) - case single :: Nil => (nodes.updated(hash, single), collisions - hash) + case single :: Nil => (single, Nil) case _ => - // lowest toString owns the position, independent of the order the nodes were added val sorted = claimants.sortBy(_.toString) - (nodes.updated(hash, sorted.head), collisions.updated(hash, sorted.tail)) + (sorted.head, sorted.tail) } // adds the virtual nodes of `node`, replacing any previous virtual nodes of the same node