This is an automated email from the ASF dual-hosted git repository.
Aias00 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/shenyu.git
The following commit(s) were added to refs/heads/master by this push:
new 0bcf9ec57c perf(loadbalancer): cache consistent hash rings (#7161)
0bcf9ec57c is described below
commit 0bcf9ec57ccf75a4ddf96f1fa429614c00530af1
Author: Liming Deng <[email protected]>
AuthorDate: Wed Sep 30 09:58:35 2026 +0800
perf(loadbalancer): cache consistent hash rings (#7161)
Co-authored-by: aias00 <[email protected]>
Co-authored-by: zhengpeng <[email protected]>
---
.../shenyu/loadbalancer/spi/HashLoadBalancer.java | 53 ++++++++++++++--------
.../loadbalancer/spi/HashLoadBalancerTest.java | 29 ++++++++++++
2 files changed, 64 insertions(+), 18 deletions(-)
diff --git
a/shenyu-loadbalancer/src/main/java/org/apache/shenyu/loadbalancer/spi/HashLoadBalancer.java
b/shenyu-loadbalancer/src/main/java/org/apache/shenyu/loadbalancer/spi/HashLoadBalancer.java
index 8de25c0e2f..06f5380d93 100644
---
a/shenyu-loadbalancer/src/main/java/org/apache/shenyu/loadbalancer/spi/HashLoadBalancer.java
+++
b/shenyu-loadbalancer/src/main/java/org/apache/shenyu/loadbalancer/spi/HashLoadBalancer.java
@@ -17,6 +17,8 @@
package org.apache.shenyu.loadbalancer.spi;
+import org.apache.shenyu.common.cache.WindowTinyLFUMap;
+import org.apache.shenyu.common.constant.Constants;
import org.apache.shenyu.loadbalancer.entity.LoadBalanceData;
import org.apache.shenyu.loadbalancer.entity.Upstream;
import org.apache.shenyu.spi.Join;
@@ -24,9 +26,11 @@ import org.apache.shenyu.spi.Join;
import java.nio.charset.StandardCharsets;
import java.security.MessageDigest;
import java.security.NoSuchAlgorithmException;
+import java.util.ArrayList;
import java.util.List;
+import java.util.Map;
import java.util.SortedMap;
-import java.util.concurrent.ConcurrentSkipListMap;
+import java.util.TreeMap;
import java.util.stream.IntStream;
/**
@@ -39,33 +43,38 @@ public class HashLoadBalancer extends AbstractLoadBalancer {
* virtual node used to solve unbalanced load.
*/
private static final int VIRTUAL_NODE_NUM = 5;
+
+ private static final ThreadLocal<MessageDigest> MD5 =
ThreadLocal.withInitial(HashLoadBalancer::newMessageDigest);
+
+ private final Map<List<String>, SortedMap<Long, Integer>> ringCache = new
WindowTinyLFUMap<>(Constants.CACHE_MAX_COUNT);
@Override
public Upstream doSelect(final List<Upstream> upstreamList, final
LoadBalanceData data) {
- final ConcurrentSkipListMap<Long, Upstream> treeMap = new
ConcurrentSkipListMap<>();
- upstreamList.forEach(upstream -> IntStream.range(0,
VIRTUAL_NODE_NUM).forEach(i -> {
- long addressHash = hash("SHENYU-" + upstream.getUrl() + "-HASH-" +
i);
- treeMap.put(addressHash, upstream);
- }));
+ final List<String> ringKey = new ArrayList<>(upstreamList.size());
+ upstreamList.forEach(upstream -> ringKey.add(upstream.getUrl()));
+ final SortedMap<Long, Integer> treeMap =
ringCache.computeIfAbsent(List.copyOf(ringKey), this::buildRing);
long hash = hash(data.getIp());
- SortedMap<Long, Upstream> lastRing = treeMap.tailMap(hash);
+ SortedMap<Long, Integer> lastRing = treeMap.tailMap(hash);
if (!lastRing.isEmpty()) {
- return lastRing.get(lastRing.firstKey());
+ return upstreamList.get(lastRing.get(lastRing.firstKey()));
}
- return treeMap.firstEntry().getValue();
+ return upstreamList.get(treeMap.get(treeMap.firstKey()));
+ }
+
+ private SortedMap<Long, Integer> buildRing(final List<String>
upstreamUrls) {
+ final SortedMap<Long, Integer> treeMap = new TreeMap<>();
+ IntStream.range(0, upstreamUrls.size()).forEach(index ->
+ IntStream.range(0, VIRTUAL_NODE_NUM).forEach(virtualNode -> {
+ long addressHash = hash("SHENYU-" +
upstreamUrls.get(index) + "-HASH-" + virtualNode);
+ treeMap.put(addressHash, index);
+ }));
+ return treeMap;
}
private static long hash(final String key) {
- // md5 byte
- MessageDigest md5;
- try {
- md5 = MessageDigest.getInstance("MD5");
- } catch (NoSuchAlgorithmException e) {
- throw new RuntimeException("MD5 not supported", e);
- }
+ MessageDigest md5 = MD5.get();
md5.reset();
- byte[] keyBytes;
- keyBytes = key.getBytes(StandardCharsets.UTF_8);
+ byte[] keyBytes = key.getBytes(StandardCharsets.UTF_8);
md5.update(keyBytes);
byte[] digest = md5.digest();
// hash code, Truncate to 32-bits
@@ -75,4 +84,12 @@ public class HashLoadBalancer extends AbstractLoadBalancer {
| (digest[0] & 0xFF);
return hashCode & 0xffffffffL;
}
+
+ private static MessageDigest newMessageDigest() {
+ try {
+ return MessageDigest.getInstance("MD5");
+ } catch (NoSuchAlgorithmException e) {
+ throw new RuntimeException("MD5 not supported", e);
+ }
+ }
}
diff --git
a/shenyu-loadbalancer/src/test/java/org/apache/shenyu/loadbalancer/spi/HashLoadBalancerTest.java
b/shenyu-loadbalancer/src/test/java/org/apache/shenyu/loadbalancer/spi/HashLoadBalancerTest.java
index a58beff641..13db6be6c9 100644
---
a/shenyu-loadbalancer/src/test/java/org/apache/shenyu/loadbalancer/spi/HashLoadBalancerTest.java
+++
b/shenyu-loadbalancer/src/test/java/org/apache/shenyu/loadbalancer/spi/HashLoadBalancerTest.java
@@ -21,10 +21,14 @@ import
org.apache.shenyu.loadbalancer.entity.LoadBalanceData;
import org.apache.shenyu.loadbalancer.entity.Upstream;
import org.junit.jupiter.api.Test;
+import java.lang.reflect.Field;
import java.util.ArrayList;
import java.util.List;
+import java.util.Map;
import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertNotSame;
+import static org.junit.jupiter.api.Assertions.assertSame;
/**
* HashLoadBalancer unit test.
@@ -43,4 +47,29 @@ class HashLoadBalancerTest {
assertEquals(upstreamList.get(2).getUrl(), upstream.getUrl());
}
+ @Test
+ void shouldReuseRingWithoutReturningStaleUpstreamInstances() throws
Exception {
+ final HashLoadBalancer hashLoadBalancer = new HashLoadBalancer();
+ final List<Upstream> first = upstreams();
+ final List<Upstream> refreshed = upstreams();
+
+ Upstream firstSelection = hashLoadBalancer.doSelect(first, new
LoadBalanceData());
+ Upstream refreshedSelection = hashLoadBalancer.doSelect(refreshed, new
LoadBalanceData());
+
+ assertEquals(first.indexOf(firstSelection),
refreshed.indexOf(refreshedSelection));
+ assertNotSame(firstSelection, refreshedSelection);
+ assertSame(refreshed.get(refreshed.indexOf(refreshedSelection)),
refreshedSelection);
+ Field ringCacheField =
HashLoadBalancer.class.getDeclaredField("ringCache");
+ ringCacheField.setAccessible(true);
+ assertEquals(1, ((Map<?, ?>)
ringCacheField.get(hashLoadBalancer)).size());
+ }
+
+ private List<Upstream> upstreams() {
+ final List<Upstream> result = new ArrayList<>();
+ result.add(Upstream.builder().url("http://1.1.1.1/api").build());
+ result.add(Upstream.builder().url("http://2.2.2.2/api").build());
+ result.add(Upstream.builder().url("http://3.3.3.3/api").build());
+ return result;
+ }
+
}