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"));
+    }
 }

Reply via email to