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 f3167af37d fix: evict stale upstream countMap entries in
LeastActiveLoadBalance (#6900)
f3167af37d is described below
commit f3167af37d7caaeeae55e18fdeb2b9cbd258f5b1
Author: Southern <[email protected]>
AuthorDate: Sat Aug 15 15:30:54 2026 +0800
fix: evict stale upstream countMap entries in LeastActiveLoadBalance
(#6900)
* fix: evict stale upstream countMap entries in LeastActiveLoadBalance
(#6891)
The countMap in LeastActiveLoadBalance accumulates an entry for every
upstream domain ever observed, but had no removal path: when an upstream was
removed from the list, its entry lingered indefinitely, leaking memory
proportionally to historical upstream churn and adding O(n) scan overhead to
every doSelect call. Now stale entries whose domains are absent from the
current upstream list are removed before the least active domain is selected.
* fix: use time-based eviction for LeastActiveLoadBalance countMap (#6891)
The previous retainAll-based cleanup removed every entry that was not
present in the current selector's upstream list, but the countMap is a
JVM-wide singleton shared by all selectors, so a cleanup triggered by one
selector deleted other selectors' live entries and lost their accumulated
counts.
Rework the cleanup to be time-driven instead: live entries have their
lastUpdate refreshed on every request, and only entries untouched for
longer than the recycle period (60s) are evicted. This bounds the memory
leak while keeping each selector's counts isolated, and the eviction is
throttled so it adds no per-request overhead. Tests are updated to cover
stale-entry eviction and cross-selector isolation.
* docs: Add inline comment explaining 60s grace window in doSelect (#6891)
Explain that the time-based eviction via removeIf tolerates a temporarily
absent domain during the 60s recycle period, with no correctness impact — only
a bounded memory tail.
---------
Co-authored-by: aias00 <[email protected]>
---
.../loadbalancer/spi/LeastActiveLoadBalance.java | 87 +++++++++++++++++++--
.../spi/LeastActiveLoadBalanceTest.java | 88 ++++++++++++++++++++++
2 files changed, 167 insertions(+), 8 deletions(-)
diff --git
a/shenyu-loadbalancer/src/main/java/org/apache/shenyu/loadbalancer/spi/LeastActiveLoadBalance.java
b/shenyu-loadbalancer/src/main/java/org/apache/shenyu/loadbalancer/spi/LeastActiveLoadBalance.java
index 5d117a1c08..aa6eb32199 100644
---
a/shenyu-loadbalancer/src/main/java/org/apache/shenyu/loadbalancer/spi/LeastActiveLoadBalance.java
+++
b/shenyu-loadbalancer/src/main/java/org/apache/shenyu/loadbalancer/spi/LeastActiveLoadBalance.java
@@ -24,8 +24,11 @@ import org.apache.shenyu.spi.Join;
import java.util.Comparator;
import java.util.List;
import java.util.Map;
-import java.util.Optional;
+import java.util.Objects;
import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.ConcurrentMap;
+import java.util.concurrent.atomic.AtomicBoolean;
+import java.util.concurrent.atomic.AtomicLong;
import java.util.stream.Collectors;
/**
@@ -34,26 +37,94 @@ import java.util.stream.Collectors;
@Join
public class LeastActiveLoadBalance extends AbstractLoadBalancer {
- private final Map<String, Long> countMap = new ConcurrentHashMap<>();
+ private final int recyclePeriod = 60000;
+
+ private final ConcurrentMap<String, ActiveCount> countMap = new
ConcurrentHashMap<>(16);
+
+ private final AtomicBoolean updateLock = new AtomicBoolean();
+
+ private volatile long lastRecycle;
@Override
protected Upstream doSelect(final List<Upstream> upstreamList, final
LoadBalanceData data) {
+ long now = System.currentTimeMillis();
Map<String, Upstream> domainMap = upstreamList.stream()
.collect(Collectors.toConcurrentMap(Upstream::buildDomain,
upstream -> upstream));
- domainMap.keySet().stream()
- .filter(key -> !countMap.containsKey(key))
- .forEach(domain -> countMap.put(domain, Long.MIN_VALUE));
+ domainMap.keySet().forEach(domain -> {
+ ActiveCount activeCount = countMap.computeIfAbsent(domain, key ->
new ActiveCount(now));
+ activeCount.setLastUpdate(now);
+ });
final String domain = countMap.entrySet().stream()
// Ensure that the filtered domain is included in the
domainMap.
.filter(entry -> domainMap.containsKey(entry.getKey()))
- .min(Comparator.comparingLong(Map.Entry::getValue))
+ .min(Comparator.comparingLong(entry ->
entry.getValue().getCount()))
.map(Map.Entry::getKey)
.orElse(upstreamList.get(0).buildDomain());
- countMap.computeIfPresent(domain, (key, activated) ->
Optional.of(activated).orElse(Long.MIN_VALUE) + 1);
+ ActiveCount activeCount = countMap.get(domain);
+ if (Objects.nonNull(activeCount)) {
+ activeCount.increase();
+ }
+
+ // A removed domain's entry lingers for up to recyclePeriod, safely
excluded from selection meanwhile.
+ if (!updateLock.get() && now - lastRecycle > recyclePeriod &&
updateLock.compareAndSet(false, true)) {
+ try {
+ countMap.entrySet().removeIf(item -> now -
item.getValue().getLastUpdate() > recyclePeriod);
+ lastRecycle = now;
+ } finally {
+ updateLock.set(false);
+ }
+ }
return domainMap.get(domain);
}
-
+
+ /**
+ * The type Active count.
+ */
+ protected static class ActiveCount {
+
+ private final AtomicLong count = new AtomicLong(Long.MIN_VALUE);
+
+ private volatile long lastUpdate;
+
+ ActiveCount(final long lastUpdate) {
+ this.lastUpdate = lastUpdate;
+ }
+
+ /**
+ * Increase count.
+ */
+ void increase() {
+ count.addAndGet(1);
+ }
+
+ /**
+ * Get count.
+ *
+ * @return the count
+ */
+ long getCount() {
+ return count.get();
+ }
+
+ /**
+ * Gets last update.
+ *
+ * @return the last update
+ */
+ long getLastUpdate() {
+ return lastUpdate;
+ }
+
+ /**
+ * Sets last update.
+ *
+ * @param lastUpdate the last update
+ */
+ void setLastUpdate(final long lastUpdate) {
+ this.lastUpdate = lastUpdate;
+ }
+ }
}
diff --git
a/shenyu-loadbalancer/src/test/java/org/apache/shenyu/loadbalancer/spi/LeastActiveLoadBalanceTest.java
b/shenyu-loadbalancer/src/test/java/org/apache/shenyu/loadbalancer/spi/LeastActiveLoadBalanceTest.java
index 44a6851dc2..55ad640b29 100644
---
a/shenyu-loadbalancer/src/test/java/org/apache/shenyu/loadbalancer/spi/LeastActiveLoadBalanceTest.java
+++
b/shenyu-loadbalancer/src/test/java/org/apache/shenyu/loadbalancer/spi/LeastActiveLoadBalanceTest.java
@@ -22,8 +22,10 @@ import org.apache.shenyu.loadbalancer.entity.Upstream;
import org.junit.jupiter.api.Assertions;
import org.junit.jupiter.api.Test;
+import java.lang.reflect.Field;
import java.util.ArrayList;
import java.util.List;
+import java.util.Map;
/**
* The type least activity load balance test.
@@ -56,4 +58,90 @@ public class LeastActiveLoadBalanceTest {
Assertions.assertTrue(upstream.getUrl().equals("baidu.com") &&
upstream1.getUrl().equals("pro.jd.com")
|| upstream1.getUrl().equals("baidu.com") &&
upstream.getUrl().equals("pro.jd.com"));
}
+
+ @Test
+ public void testRemoveStaleCountMapEntries() throws Exception {
+ buildUpstreamList();
+ final LeastActiveLoadBalance leastActiveLoadBalance = new
LeastActiveLoadBalance();
+ leastActiveLoadBalance.doSelect(onlyOneList, new LoadBalanceData());
+ Assertions.assertEquals(2, getCountMap(leastActiveLoadBalance).size());
+ onlyOneList.remove(1);
+ // The periodic cleanup is time-driven, simulate the elapsed time so
that the entry of the
+ // removed upstream becomes stale and is evicted by the next cleanup.
+ long old = System.currentTimeMillis() - 60_001;
+ setLastUpdate(leastActiveLoadBalance, "https://pro.jd.com", old);
+ setLastRecycle(leastActiveLoadBalance, old);
+ leastActiveLoadBalance.doSelect(onlyOneList, new LoadBalanceData());
+ Map<String, LeastActiveLoadBalance.ActiveCount> countMap =
getCountMap(leastActiveLoadBalance);
+ Assertions.assertEquals(1, countMap.size());
+ Assertions.assertTrue(countMap.containsKey("https://baidu.com"));
+ Assertions.assertFalse(countMap.containsKey("https://pro.jd.com"));
+ }
+
+ @Test
+ public void testCountMapNotAffectOtherSelector() throws Exception {
+ buildUpstreamList();
+ final List<Upstream> anotherList = new ArrayList<>();
+
anotherList.add(Upstream.builder().url("jd.com").protocol("https://").build());
+ final LeastActiveLoadBalance leastActiveLoadBalance = new
LeastActiveLoadBalance();
+ leastActiveLoadBalance.doSelect(onlyOneList, new LoadBalanceData());
+ leastActiveLoadBalance.doSelect(anotherList, new LoadBalanceData());
+ Assertions.assertEquals(3, getCountMap(leastActiveLoadBalance).size());
+ // Selector 1 removes pro.jd.com, and the cleanup must not evict the
live entry jd.com
+ // which only belongs to selector 2.
+ onlyOneList.remove(1);
+ long old = System.currentTimeMillis() - 60_001;
+ setLastUpdate(leastActiveLoadBalance, "https://pro.jd.com", old);
+ setLastRecycle(leastActiveLoadBalance, old);
+ leastActiveLoadBalance.doSelect(onlyOneList, new LoadBalanceData());
+ Map<String, LeastActiveLoadBalance.ActiveCount> countMap =
getCountMap(leastActiveLoadBalance);
+ Assertions.assertEquals(2, countMap.size());
+ Assertions.assertTrue(countMap.containsKey("https://baidu.com"));
+ Assertions.assertTrue(countMap.containsKey("https://jd.com"));
+ Assertions.assertFalse(countMap.containsKey("https://pro.jd.com"));
+ }
+
+ @Test
+ public void testCountMapNotAffectOtherSelectorWhenFirstRemoved() throws
Exception {
+ buildUpstreamList();
+ final List<Upstream> anotherList = new ArrayList<>();
+
anotherList.add(Upstream.builder().url("jd.com").protocol("https://").build());
+ final LeastActiveLoadBalance leastActiveLoadBalance = new
LeastActiveLoadBalance();
+ leastActiveLoadBalance.doSelect(onlyOneList, new LoadBalanceData());
+ leastActiveLoadBalance.doSelect(anotherList, new LoadBalanceData());
+ Assertions.assertEquals(3, getCountMap(leastActiveLoadBalance).size());
+ // Symmetric case: selector 1 removes the first upstream baidu.com.
The cleanup must evict
+ // baidu.com but keep pro.jd.com (still served by selector 1) and
jd.com (served by selector 2).
+ onlyOneList.remove(0);
+ long old = System.currentTimeMillis() - 60_001;
+ setLastUpdate(leastActiveLoadBalance, "https://baidu.com", old);
+ setLastRecycle(leastActiveLoadBalance, old);
+ leastActiveLoadBalance.doSelect(onlyOneList, new LoadBalanceData());
+ Map<String, LeastActiveLoadBalance.ActiveCount> countMap =
getCountMap(leastActiveLoadBalance);
+ Assertions.assertEquals(2, countMap.size());
+ Assertions.assertTrue(countMap.containsKey("https://pro.jd.com"));
+ Assertions.assertTrue(countMap.containsKey("https://jd.com"));
+ Assertions.assertFalse(countMap.containsKey("https://baidu.com"));
+ }
+
+ @SuppressWarnings("unchecked")
+ private Map<String, LeastActiveLoadBalance.ActiveCount> getCountMap(final
LeastActiveLoadBalance loadBalance) throws Exception {
+ Field field =
LeastActiveLoadBalance.class.getDeclaredField("countMap");
+ field.setAccessible(true);
+ return (Map<String, LeastActiveLoadBalance.ActiveCount>)
field.get(loadBalance);
+ }
+
+ private void setLastUpdate(final LeastActiveLoadBalance loadBalance, final
String domain, final long lastUpdate) throws Exception {
+ LeastActiveLoadBalance.ActiveCount activeCount =
getCountMap(loadBalance).get(domain);
+ Assertions.assertNotNull(activeCount);
+ Field field =
LeastActiveLoadBalance.ActiveCount.class.getDeclaredField("lastUpdate");
+ field.setAccessible(true);
+ field.set(activeCount, lastUpdate);
+ }
+
+ private void setLastRecycle(final LeastActiveLoadBalance loadBalance,
final long lastRecycle) throws Exception {
+ Field field =
LeastActiveLoadBalance.class.getDeclaredField("lastRecycle");
+ field.setAccessible(true);
+ field.set(loadBalance, lastRecycle);
+ }
}