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 cf3721cf21 Fixes #6796: Sort cached selector and rule batches once per
key (#7345)
cf3721cf21 is described below
commit cf3721cf219124401817ec681a04f65476afd45d
Author: BobSong <[email protected]>
AuthorDate: Thu Oct 1 07:36:20 2026 +0800
Fixes #6796: Sort cached selector and rule batches once per key (#7345)
---
.../shenyu/plugin/base/cache/BaseDataCache.java | 66 ++++++++++++++--------
.../plugin/base/cache/BaseDataCacheTest.java | 13 +++++
2 files changed, 54 insertions(+), 25 deletions(-)
diff --git
a/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/BaseDataCache.java
b/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/BaseDataCache.java
index fbeb778e88..745a9bdec0 100644
---
a/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/BaseDataCache.java
+++
b/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/BaseDataCache.java
@@ -25,9 +25,12 @@ import org.apache.shenyu.common.dto.SelectorData;
import java.util.ArrayList;
import java.util.Comparator;
import java.util.List;
+import java.util.Map;
import java.util.Objects;
import java.util.Optional;
+import java.util.Set;
import java.util.concurrent.ConcurrentMap;
+import java.util.function.Function;
import java.util.stream.Collectors;
/**
@@ -271,18 +274,8 @@ public final class BaseDataCache {
* @param data the rule data
*/
private void ruleAccept(final RuleData data) {
- ruleAccept(ruleMap, data);
- }
-
- private void ruleAccept(final ConcurrentMap<String, List<RuleData>>
target, final RuleData data) {
- String selectorId = data.getSelectorId();
- target.compute(selectorId, (key, value) -> {
- final List<RuleData> result = Objects.isNull(value) ? new
ArrayList<>() : new ArrayList<>(value);
- result.removeIf(rule -> Objects.equals(rule.getId(),
data.getId()));
- result.add(data);
- result.sort(Comparator.comparing(RuleData::getSort));
- return List.copyOf(result);
- });
+ ruleMap.compute(data.getSelectorId(), (key, value) ->
+ upsertSorted(value, List.of(data), RuleData::getId,
RuleData::getSort));
}
/**
@@ -291,18 +284,31 @@ public final class BaseDataCache {
* @param data the selector data
*/
private void selectorAccept(final SelectorData data) {
- selectorAccept(selectorMap, data);
+ selectorMap.compute(data.getPluginName(), (key, value) ->
+ upsertSorted(value, List.of(data), SelectorData::getId,
SelectorData::getSort));
}
- private void selectorAccept(final ConcurrentMap<String,
List<SelectorData>> target, final SelectorData data) {
- String key = data.getPluginName();
- target.compute(key, (pluginName, value) -> {
- final List<SelectorData> result = Objects.isNull(value) ? new
ArrayList<>() : new ArrayList<>(value);
- result.removeIf(selector -> Objects.equals(selector.getId(),
data.getId()));
- result.add(data);
- result.sort(Comparator.comparing(SelectorData::getSort));
- return List.copyOf(result);
- });
+ /**
+ * Merge the batch into the current list as an upsert (same-id entries are
replaced)
+ * and return a new immutable snapshot sorted by the sort key. The list is
sorted once
+ * per merge, so a batch refresh of arbitrary size costs a single sort per
key instead
+ * of one sort per element.
+ *
+ * @param current the currently cached list, may be {@code null}
+ * @param batch the incoming entries, must not contain {@code null}
elements
+ * @param idOf the id extractor used to replace entries
+ * @param sortOf the sort key extractor
+ * @param <T> the entry type
+ * @return a new immutable list sorted by the sort key
+ */
+ private <T> List<T> upsertSorted(final List<T> current, final List<T>
batch,
+ final Function<T, String> idOf, final
Function<T, Integer> sortOf) {
+ final List<T> result = Objects.isNull(current) ? new ArrayList<>() :
new ArrayList<>(current);
+ final Set<String> replacedIds =
batch.stream().map(idOf).collect(Collectors.toSet());
+ result.removeIf(item -> replacedIds.contains(idOf.apply(item)));
+ result.addAll(batch);
+ result.sort(Comparator.comparing(sortOf));
+ return List.copyOf(result);
}
/**
@@ -324,6 +330,7 @@ public final class BaseDataCache {
/**
* Merge a batch without exposing partially refreshed data to readers.
* Missing entries are retained because refresh messages may cover only
one plugin.
+ * The batch is grouped per plugin, so each plugin's list is merged and
sorted once.
*
* @param dataList the received data
*/
@@ -333,13 +340,18 @@ public final class BaseDataCache {
}
ConcurrentMap<String, List<SelectorData>> next =
Maps.newConcurrentMap();
next.putAll(selectorMap);
- dataList.forEach(data -> selectorAccept(next, data));
+ Map<String, List<SelectorData>> grouped = dataList.stream()
+ .filter(Objects::nonNull)
+ .collect(Collectors.groupingBy(SelectorData::getPluginName));
+ grouped.forEach((pluginName, batch) -> next.compute(pluginName, (key,
value) ->
+ upsertSorted(value, batch, SelectorData::getId,
SelectorData::getSort)));
selectorMap = next;
}
/**
* Merge a batch without exposing partially refreshed data to readers.
- * Missing entries are retained because refresh messages may cover only
one plugin.
+ * Missing entries are retained because refresh messages may cover only
one selector.
+ * The batch is grouped per selector, so each selector's list is merged
and sorted once.
*
* @param dataList the received data
*/
@@ -349,7 +361,11 @@ public final class BaseDataCache {
}
ConcurrentMap<String, List<RuleData>> next = Maps.newConcurrentMap();
next.putAll(ruleMap);
- dataList.forEach(data -> ruleAccept(next, data));
+ Map<String, List<RuleData>> grouped = dataList.stream()
+ .filter(Objects::nonNull)
+ .collect(Collectors.groupingBy(RuleData::getSelectorId));
+ grouped.forEach((selectorId, batch) -> next.compute(selectorId, (key,
value) ->
+ upsertSorted(value, batch, RuleData::getId,
RuleData::getSort)));
ruleMap = next;
}
}
diff --git
a/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/BaseDataCacheTest.java
b/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/BaseDataCacheTest.java
index fc2863bf71..918a4d44b2 100644
---
a/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/BaseDataCacheTest.java
+++
b/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/BaseDataCacheTest.java
@@ -467,4 +467,17 @@ public final class BaseDataCacheTest {
private RuleData rule(final String id, final int sort) {
return
RuleData.builder().id(id).pluginName("divide").selectorId("selector").sort(sort).build();
}
+
+ @Test
+ public void batchRefreshUpsertsAndSortsTheMergedListOncePerKey() {
+ cache.cacheSelectData(selector("a", 4));
+ cache.cacheRuleData(rule("a", 4));
+ cache.refreshSelectorData(List.of(selector("d", 1), selector("c", 3),
selector("a", 2), selector("e", 5)));
+ cache.refreshRuleData(List.of(rule("d", 1), rule("c", 3), rule("a",
2), rule("e", 5)));
+ // same-id entries are replaced by the batch and the merged list is
sorted by sort
+ assertEquals(List.of(selector("d", 1), selector("a", 2), selector("c",
3), selector("e", 5)),
+ cache.obtainSelectorData("divide"));
+ assertEquals(List.of(rule("d", 1), rule("a", 2), rule("c", 3),
rule("e", 5)),
+ cache.obtainRuleData("selector"));
+ }
}