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 c2cfd8c505 fix: avoid transient cache gaps during websocket refresh 
(#7035)
c2cfd8c505 is described below

commit c2cfd8c50564da08a9e6830e1ad5fc26eac764fb
Author: TheoLi0905 <[email protected]>
AuthorDate: Sun Sep 27 19:33:12 2026 +0800

    fix: avoid transient cache gaps during websocket refresh (#7035)
    
    Co-authored-by: aias00 <[email protected]>
---
 .../shenyu/plugin/base/cache/BaseDataCache.java    |  98 ++++++++---
 .../base/cache/CommonPluginDataSubscriber.java     | 114 +++++++++----
 .../plugin/base/cache/BaseDataCacheTest.java       | 182 ++++++++++++++++++++-
 .../base/cache/CommonPluginDataSubscriberTest.java | 176 ++++++++++++++++++++
 .../shenyu/sync/data/api/PluginDataSubscriber.java |  30 ++++
 .../sync/data/api/PluginDataSubscriberTest.java    |  81 +++++++++
 .../data/websocket/handler/PluginDataHandler.java  |   3 +-
 .../data/websocket/handler/RuleDataHandler.java    |   3 +-
 .../websocket/handler/SelectorDataHandler.java     |   3 +-
 .../client/ShenyuWebsocketClientTest.java          |   8 +-
 .../websocket/handler/PluginDataHandlerTest.java   |  24 ++-
 .../websocket/handler/RuleDataHandlerTest.java     |  24 ++-
 .../websocket/handler/SelectorDataHandlerTest.java |  25 ++-
 .../handler/WebsocketDataHandlerTest.java          |   4 +-
 14 files changed, 701 insertions(+), 74 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 216bd4dba8..fbeb778e88 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
@@ -40,17 +40,17 @@ public final class BaseDataCache {
     /**
      * pluginName -> PluginData.
      */
-    private static final ConcurrentMap<String, PluginData> PLUGIN_MAP = 
Maps.newConcurrentMap();
+    private static volatile ConcurrentMap<String, PluginData> pluginMap = 
Maps.newConcurrentMap();
 
     /**
      * pluginName -> SelectorData.
      */
-    private static final ConcurrentMap<String, List<SelectorData>> 
SELECTOR_MAP = Maps.newConcurrentMap();
+    private static volatile ConcurrentMap<String, List<SelectorData>> 
selectorMap = Maps.newConcurrentMap();
 
     /**
      * selectorId -> RuleData.
      */
-    private static final ConcurrentMap<String, List<RuleData>> RULE_MAP = 
Maps.newConcurrentMap();
+    private static volatile ConcurrentMap<String, List<RuleData>> ruleMap = 
Maps.newConcurrentMap();
 
     private BaseDataCache() {
     }
@@ -70,7 +70,7 @@ public final class BaseDataCache {
      * @param pluginData the plugin data
      */
     public void cachePluginData(final PluginData pluginData) {
-        Optional.ofNullable(pluginData).ifPresent(data -> 
PLUGIN_MAP.put(data.getName(), data));
+        Optional.ofNullable(pluginData).ifPresent(data -> 
pluginMap.put(data.getName(), data));
     }
     
     /**
@@ -79,7 +79,7 @@ public final class BaseDataCache {
      * @param pluginData the plugin data
      */
     public void removePluginData(final PluginData pluginData) {
-        Optional.ofNullable(pluginData).ifPresent(data -> 
PLUGIN_MAP.remove(data.getName()));
+        Optional.ofNullable(pluginData).ifPresent(data -> 
pluginMap.remove(data.getName()));
     }
     
     /**
@@ -88,14 +88,14 @@ public final class BaseDataCache {
      * @param pluginName the plugin name
      */
     public void removePluginDataByPluginName(final String pluginName) {
-        PLUGIN_MAP.remove(pluginName);
+        pluginMap.remove(pluginName);
     }
     
     /**
      * Clean plugin data.
      */
     public void cleanPluginData() {
-        PLUGIN_MAP.clear();
+        pluginMap.clear();
     }
     
     /**
@@ -114,7 +114,7 @@ public final class BaseDataCache {
      * @return the plugin data
      */
     public PluginData obtainPluginData(final String pluginName) {
-        return PLUGIN_MAP.get(pluginName);
+        return pluginMap.get(pluginName);
     }
     
     /**
@@ -133,7 +133,7 @@ public final class BaseDataCache {
      */
     public void removeSelectData(final SelectorData selectorData) {
         Optional.ofNullable(selectorData).ifPresent(data -> {
-            SELECTOR_MAP.computeIfPresent(data.getPluginName(), (key, value) 
-> {
+            selectorMap.computeIfPresent(data.getPluginName(), (key, value) -> 
{
                 final List<SelectorData> result = value.stream()
                         .filter(selector -> !Objects.equals(selector.getId(), 
data.getId()))
                         .collect(Collectors.toList());
@@ -148,14 +148,14 @@ public final class BaseDataCache {
      * @param pluginName the plugin name
      */
     public void removeSelectDataByPluginName(final String pluginName) {
-        SELECTOR_MAP.remove(pluginName);
+        selectorMap.remove(pluginName);
     }
     
     /**
      * Clean selector data.
      */
     public void cleanSelectorData() {
-        SELECTOR_MAP.clear();
+        selectorMap.clear();
     }
     
     /**
@@ -174,7 +174,7 @@ public final class BaseDataCache {
      * @return the immutable snapshot, or {@code null} if no selector data 
exists
      */
     public List<SelectorData> obtainSelectorData(final String pluginName) {
-        return SELECTOR_MAP.get(pluginName);
+        return selectorMap.get(pluginName);
     }
     
     /**
@@ -193,7 +193,7 @@ public final class BaseDataCache {
      */
     public void removeRuleData(final RuleData ruleData) {
         Optional.ofNullable(ruleData).ifPresent(data -> {
-            RULE_MAP.computeIfPresent(data.getSelectorId(), (key, value) -> {
+            ruleMap.computeIfPresent(data.getSelectorId(), (key, value) -> {
                 final List<RuleData> result = value.stream()
                         .filter(rule -> !Objects.equals(rule.getId(), 
data.getId()))
                         .collect(Collectors.toList());
@@ -208,14 +208,14 @@ public final class BaseDataCache {
      * @param selectorId the selector id
      */
     public void removeRuleDataBySelectorId(final String selectorId) {
-        RULE_MAP.remove(selectorId);
+        ruleMap.remove(selectorId);
     }
     
     /**
      * Clean rule data.
      */
     public void cleanRuleData() {
-        RULE_MAP.clear();
+        ruleMap.clear();
     }
     
     /**
@@ -234,7 +234,7 @@ public final class BaseDataCache {
      * @return the immutable snapshot, or {@code null} if no rule data exists
      */
     public List<RuleData> obtainRuleData(final String selectorId) {
-        return RULE_MAP.get(selectorId);
+        return ruleMap.get(selectorId);
     }
     
     /**
@@ -243,7 +243,7 @@ public final class BaseDataCache {
      * @return the plugin map
      */
     public ConcurrentMap<String, PluginData> getPluginMap() {
-        return PLUGIN_MAP;
+        return pluginMap;
     }
     
     /**
@@ -252,7 +252,7 @@ public final class BaseDataCache {
      * @return the selector map
      */
     public ConcurrentMap<String, List<SelectorData>> getSelectorMap() {
-        return SELECTOR_MAP;
+        return selectorMap;
     }
     
     /**
@@ -261,7 +261,7 @@ public final class BaseDataCache {
      * @return the rule map
      */
     public ConcurrentMap<String, List<RuleData>> getRuleMap() {
-        return RULE_MAP;
+        return ruleMap;
     }
     
 
@@ -271,8 +271,12 @@ 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();
-        RULE_MAP.compute(selectorId, (key, value) -> {
+        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);
@@ -287,8 +291,12 @@ public final class BaseDataCache {
      * @param data the selector data
      */
     private void selectorAccept(final SelectorData data) {
+        selectorAccept(selectorMap, data);
+    }
+
+    private void selectorAccept(final ConcurrentMap<String, 
List<SelectorData>> target, final SelectorData data) {
         String key = data.getPluginName();
-        SELECTOR_MAP.compute(key, (pluginName, value) -> {
+        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);
@@ -296,4 +304,52 @@ public final class BaseDataCache {
             return List.copyOf(result);
         });
     }
+
+    /**
+     * Merge a batch without exposing partially refreshed data to readers.
+     * Missing entries are retained because refresh messages may cover only 
one plugin.
+     *
+     * @param dataList the received data
+     */
+    void refreshPluginData(final List<PluginData> dataList) {
+        if (dataList.isEmpty()) {
+            return;
+        }
+        ConcurrentMap<String, PluginData> next = Maps.newConcurrentMap();
+        next.putAll(pluginMap);
+        dataList.forEach(data -> next.put(data.getName(), data));
+        pluginMap = next;
+    }
+
+    /**
+     * Merge a batch without exposing partially refreshed data to readers.
+     * Missing entries are retained because refresh messages may cover only 
one plugin.
+     *
+     * @param dataList the received data
+     */
+    void refreshSelectorData(final List<SelectorData> dataList) {
+        if (dataList.isEmpty()) {
+            return;
+        }
+        ConcurrentMap<String, List<SelectorData>> next = 
Maps.newConcurrentMap();
+        next.putAll(selectorMap);
+        dataList.forEach(data -> selectorAccept(next, data));
+        selectorMap = next;
+    }
+
+    /**
+     * Merge a batch without exposing partially refreshed data to readers.
+     * Missing entries are retained because refresh messages may cover only 
one plugin.
+     *
+     * @param dataList the received data
+     */
+    void refreshRuleData(final List<RuleData> dataList) {
+        if (dataList.isEmpty()) {
+            return;
+        }
+        ConcurrentMap<String, List<RuleData>> next = Maps.newConcurrentMap();
+        next.putAll(ruleMap);
+        dataList.forEach(data -> ruleAccept(next, data));
+        ruleMap = next;
+    }
 }
diff --git 
a/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriber.java
 
b/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriber.java
index d71691056a..9f2d6f84f9 100644
--- 
a/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriber.java
+++ 
b/shenyu-plugin/shenyu-plugin-base/src/main/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriber.java
@@ -214,44 +214,15 @@ public class CommonPluginDataSubscriber implements 
PluginDataSubscriber {
                     .ifPresent(handler -> handler.handlerPlugin(pluginData));
 
             BaseDataCache.getInstance().cachePluginData(pluginData);
-            // update enabled plugins
-            PluginHandlerEventEnum state = 
Boolean.TRUE.equals(pluginData.getEnabled())
-                    ? PluginHandlerEventEnum.ENABLED : 
PluginHandlerEventEnum.DISABLED;
-            eventPublisher.publishEvent(new PluginHandlerEvent(state, 
pluginData));
-            // sorted plugin
-            sortPluginIfOrderChange(oldPluginData, pluginData);
-            
-            final String pluginName = pluginData.getName();
-            // if update plugin, remove selector and rule match cache/trie 
cache
-            if (selectorMatchConfig.getCache().getEnabled()) {
-                MatchDataCache.getInstance().removeSelectorData(pluginName);
-            }
-            if (ruleMatchCacheConfig.getCache().getEnabled()) {
-                MatchDataCache.getInstance().removeRuleData(pluginName);
-            }
+            notifyPluginData(oldPluginData, pluginData);
         } else if (data instanceof SelectorData) {
             SelectorData selectorData = (SelectorData) data;
             BaseDataCache.getInstance().cacheSelectData(selectorData);
-            Optional.ofNullable(handlerMap.get(selectorData.getPluginName()))
-                    .ifPresent(handler -> 
handler.handlerSelector(selectorData));
-            // remove match cache
-            if (selectorMatchConfig.getCache().getEnabled()) {
-                
MatchDataCache.getInstance().removeSelectorData(selectorData.getPluginName(), 
selectorData.getId());
-                
MatchDataCache.getInstance().removeEmptySelectorData(selectorData.getPluginName());
-            }
-            if (ruleMatchCacheConfig.getCache().getEnabled()) {
-                
MatchDataCache.getInstance().removeRuleDataBySelector(selectorData.getPluginName(),
 selectorData.getId());
-                
MatchDataCache.getInstance().removeEmptyRuleData(selectorData.getPluginName());
-            }
+            handleSelectorData(selectorData);
         } else if (data instanceof RuleData) {
             RuleData ruleData = (RuleData) data;
             BaseDataCache.getInstance().cacheRuleData(ruleData);
-            Optional.ofNullable(handlerMap.get(ruleData.getPluginName()))
-                    .ifPresent(handler -> handler.handlerRule(ruleData));
-            if (ruleMatchCacheConfig.getCache().getEnabled()) {
-                
MatchDataCache.getInstance().removeRuleData(ruleData.getPluginName(), 
ruleData.getId());
-                
MatchDataCache.getInstance().removeEmptyRuleData(ruleData.getPluginName());
-            }
+            handleRuleData(ruleData);
         }
     }
 
@@ -310,4 +281,83 @@ public class CommonPluginDataSubscriber implements 
PluginDataSubscriber {
         }
     }
 
+    private void notifyPluginData(final PluginData oldPluginData, final 
PluginData pluginData) {
+        // update enabled plugins
+        PluginHandlerEventEnum state = 
Boolean.TRUE.equals(pluginData.getEnabled())
+                ? PluginHandlerEventEnum.ENABLED : 
PluginHandlerEventEnum.DISABLED;
+        eventPublisher.publishEvent(new PluginHandlerEvent(state, pluginData));
+        // sorted plugin
+        sortPluginIfOrderChange(oldPluginData, pluginData);
+
+        final String pluginName = pluginData.getName();
+        // if update plugin, remove selector and rule match cache/trie cache
+        if (selectorMatchConfig.getCache().getEnabled()) {
+            MatchDataCache.getInstance().removeSelectorData(pluginName);
+        }
+        if (ruleMatchCacheConfig.getCache().getEnabled()) {
+            MatchDataCache.getInstance().removeRuleData(pluginName);
+        }
+    }
+
+    private void handleSelectorData(final SelectorData selectorData) {
+        Optional.ofNullable(handlerMap.get(selectorData.getPluginName()))
+                .ifPresent(handler -> handler.handlerSelector(selectorData));
+        invalidateSelectorMatchCache(selectorData);
+    }
+
+    private void invalidateSelectorMatchCache(final SelectorData selectorData) 
{
+        // remove match cache
+        if (selectorMatchConfig.getCache().getEnabled()) {
+            
MatchDataCache.getInstance().removeSelectorData(selectorData.getPluginName(), 
selectorData.getId());
+            
MatchDataCache.getInstance().removeEmptySelectorData(selectorData.getPluginName());
+        }
+        if (ruleMatchCacheConfig.getCache().getEnabled()) {
+            
MatchDataCache.getInstance().removeRuleDataBySelector(selectorData.getPluginName(),
 selectorData.getId());
+            
MatchDataCache.getInstance().removeEmptyRuleData(selectorData.getPluginName());
+        }
+    }
+
+    private void handleRuleData(final RuleData ruleData) {
+        Optional.ofNullable(handlerMap.get(ruleData.getPluginName()))
+                .ifPresent(handler -> handler.handlerRule(ruleData));
+        invalidateRuleMatchCache(ruleData);
+    }
+
+    private void invalidateRuleMatchCache(final RuleData ruleData) {
+        if (ruleMatchCacheConfig.getCache().getEnabled()) {
+            
MatchDataCache.getInstance().removeRuleData(ruleData.getPluginName(), 
ruleData.getId());
+            
MatchDataCache.getInstance().removeEmptyRuleData(ruleData.getPluginName());
+        }
+    }
+
+    @Override
+    public void onPluginRefresh(final List<PluginData> dataList) {
+        if (CollectionUtils.isEmpty(dataList)) {
+            return;
+        }
+        dataList.forEach(data -> 
Optional.ofNullable(handlerMap.get(data.getName()))
+                .ifPresent(handler -> handler.handlerPlugin(data)));
+        BaseDataCache.getInstance().refreshPluginData(dataList);
+        // Legacy refresh removed the old entries before subscription, so it 
always notified sorting.
+        dataList.forEach(data -> notifyPluginData(null, data));
+    }
+
+    @Override
+    public void onSelectorRefresh(final List<SelectorData> dataList) {
+        if (CollectionUtils.isEmpty(dataList)) {
+            return;
+        }
+        BaseDataCache.getInstance().refreshSelectorData(dataList);
+        dataList.forEach(this::handleSelectorData);
+    }
+
+    @Override
+    public void onRuleRefresh(final List<RuleData> dataList) {
+        if (CollectionUtils.isEmpty(dataList)) {
+            return;
+        }
+        BaseDataCache.getInstance().refreshRuleData(dataList);
+        dataList.forEach(this::handleRuleData);
+    }
+
 }
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 e2c0c5fc5a..fc2863bf71 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
@@ -21,18 +21,30 @@ import com.google.common.collect.Lists;
 import org.apache.shenyu.common.dto.PluginData;
 import org.apache.shenyu.common.dto.RuleData;
 import org.apache.shenyu.common.dto.SelectorData;
+import org.junit.jupiter.api.AfterEach;
+import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 
 import java.lang.reflect.Field;
+import java.util.AbstractList;
 import java.util.Iterator;
 import java.util.List;
 import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.CountDownLatch;
+import java.util.concurrent.ExecutorService;
+import java.util.concurrent.Executors;
+import java.util.concurrent.Future;
+import java.util.concurrent.TimeUnit;
 
 import static org.junit.jupiter.api.Assertions.assertDoesNotThrow;
 import static org.junit.jupiter.api.Assertions.assertEquals;
+import static org.junit.jupiter.api.Assertions.assertFalse;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNotSame;
 import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
 import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.junit.jupiter.api.Assertions.assertTrue;
 
 /**
  * Test cases for BaseDataCache.
@@ -40,11 +52,17 @@ import static org.junit.jupiter.api.Assertions.assertThrows;
 @SuppressWarnings("unchecked")
 public final class BaseDataCacheTest {
 
-    private final String pluginMapStr = "PLUGIN_MAP";
+    private final BaseDataCache cache = BaseDataCache.getInstance();
 
-    private final String selectorMapStr = "SELECTOR_MAP";
+    private final CountDownLatch preparing = new CountDownLatch(1);
 
-    private final String ruleMapStr = "RULE_MAP";
+    private final CountDownLatch publish = new CountDownLatch(1);
+
+    private final String pluginMapStr = "pluginMap";
+
+    private final String selectorMapStr = "selectorMap";
+
+    private final String ruleMapStr = "ruleMap";
 
     private final String mockName1 = "MOCK_NAME_1";
     
@@ -291,4 +309,162 @@ public final class BaseDataCacheTest {
         pluginMapField.setAccessible(true);
         return (ConcurrentHashMap) pluginMapField.get(baseDataCache);
     }
+
+    @BeforeEach
+    @AfterEach
+    public void clearCaches() {
+        cache.cleanPluginData();
+        cache.cleanSelectorData();
+        cache.cleanRuleData();
+    }
+
+    @Test
+    public void pluginCandidateConstructionKeepsOldMapVisible() throws 
Exception {
+        final PluginData old = plugin("divide", 1);
+        final PluginData updated = plugin("divide", 2);
+        final PluginData added = plugin("rewrite", 3);
+        cache.cachePluginData(old);
+        final var previous = cache.getPluginMap();
+        assertAtomicRefresh(() -> 
cache.refreshPluginData(pauseBeforeLast(List.of(updated, added))), () -> {
+            assertSame(previous, cache.getPluginMap());
+            assertSame(old, cache.obtainPluginData("divide"));
+            assertNull(cache.obtainPluginData("rewrite"));
+        });
+        assertSame(updated, cache.obtainPluginData("divide"));
+        assertSame(added, cache.obtainPluginData("rewrite"));
+    }
+
+    @Test
+    public void selectorRefreshRetainsMissingDataAndPublishesSortedLists() 
throws Exception {
+        final SelectorData a = selector("a", 1);
+        final SelectorData old = selector("b", 2);
+        final SelectorData updated = selector("b", 4);
+        final SelectorData added = selector("c", 3);
+        cache.cacheSelectData(a);
+        cache.cacheSelectData(old);
+        final var previous = cache.getSelectorMap();
+        final var previousList = cache.obtainSelectorData("divide");
+        assertAtomicRefresh(() -> 
cache.refreshSelectorData(pauseBeforeLast(List.of(updated, added))), () -> {
+            assertSame(previous, cache.getSelectorMap());
+            assertEquals(List.of(a, old), cache.obtainSelectorData("divide"));
+        });
+        assertNotSame(previous, cache.getSelectorMap());
+        assertEquals(List.of(a, added, updated), 
cache.obtainSelectorData("divide"));
+        assertEquals(List.of(a, old), previousList);
+        assertThrows(UnsupportedOperationException.class, () -> 
cache.obtainSelectorData("divide").clear());
+    }
+
+    @Test
+    public void ruleRefreshRetainsMissingDataAndPublishesSortedLists() throws 
Exception {
+        final RuleData a = rule("a", 1);
+        final RuleData old = rule("b", 2);
+        final RuleData updated = rule("b", 4);
+        final RuleData added = rule("c", 3);
+        cache.cacheRuleData(a);
+        cache.cacheRuleData(old);
+        final var previous = cache.getRuleMap();
+        final var previousList = cache.obtainRuleData("selector");
+        assertAtomicRefresh(() -> 
cache.refreshRuleData(pauseBeforeLast(List.of(updated, added))), () -> {
+            assertSame(previous, cache.getRuleMap());
+            assertEquals(List.of(a, old), cache.obtainRuleData("selector"));
+        });
+        assertNotSame(previous, cache.getRuleMap());
+        assertEquals(List.of(a, added, updated), 
cache.obtainRuleData("selector"));
+        assertEquals(List.of(a, old), previousList);
+        assertThrows(UnsupportedOperationException.class, () -> 
cache.obtainRuleData("selector").clear());
+    }
+
+    @Test
+    public void invalidCandidatesDoNotPublish() {
+        cache.cachePluginData(plugin("divide", 1));
+        cache.cacheSelectData(selector("a", 1));
+        cache.cacheRuleData(rule("a", 1));
+        final var plugins = cache.getPluginMap();
+        final var selectors = cache.getSelectorMap();
+        final var rules = cache.getRuleMap();
+        assertThrows(NullPointerException.class, () -> 
cache.refreshPluginData(List.of(plugin("divide", 2), new PluginData())));
+        assertThrows(NullPointerException.class, () -> 
cache.refreshSelectorData(List.of(selector("b", 2), new SelectorData())));
+        assertThrows(NullPointerException.class, () -> 
cache.refreshRuleData(List.of(rule("b", 2), new RuleData())));
+        assertSame(plugins, cache.getPluginMap());
+        assertSame(selectors, cache.getSelectorMap());
+        assertSame(rules, cache.getRuleMap());
+    }
+
+    @Test
+    public void batchesKeepUnrelatedGroupsAndIncrementalListsImmutable() {
+        final SelectorData otherSelector = 
SelectorData.builder().id("other").pluginName("jwt").sort(1).build();
+        final RuleData otherRule = 
RuleData.builder().id("other").selectorId("other").pluginName("jwt").sort(1).build();
+        cache.cacheSelectData(otherSelector);
+        cache.cacheRuleData(otherRule);
+        cache.refreshSelectorData(List.of(selector("b", 2), selector("a", 1)));
+        cache.refreshRuleData(List.of(rule("b", 2), rule("a", 1)));
+        assertEquals(List.of(otherSelector), cache.obtainSelectorData("jwt"));
+        assertEquals(List.of(otherRule), cache.obtainRuleData("other"));
+        final var selectors = cache.obtainSelectorData("divide");
+        final var rules = cache.obtainRuleData("selector");
+        cache.cacheSelectData(selector("a", 3));
+        cache.cacheRuleData(rule("a", 3));
+        cache.removeSelectData(selector("b", 2));
+        cache.removeRuleData(rule("b", 2));
+        assertEquals(List.of(selector("a", 1), selector("b", 2)), selectors);
+        assertEquals(List.of(rule("a", 1), rule("b", 2)), rules);
+    }
+
+    private <T> List<T> pauseBeforeLast(final List<T> data) {
+        return new AbstractList<>() {
+            @Override
+            public T get(final int index) {
+                if (index == data.size() - 1) {
+                    try {
+                        pausePreparation();
+                    } catch (InterruptedException e) {
+                        Thread.currentThread().interrupt();
+                        throw new IllegalStateException(e);
+                    }
+                }
+                return data.get(index);
+            }
+
+            @Override
+            public int size() {
+                return data.size();
+            }
+        };
+    }
+
+    private Object pausePreparation() throws InterruptedException {
+        preparing.countDown();
+        assertTrue(publish.await(10, TimeUnit.SECONDS));
+        return null;
+    }
+
+    private void assertAtomicRefresh(final Runnable refresh, final Runnable 
readOld) throws Exception {
+        final ExecutorService executor = Executors.newSingleThreadExecutor();
+        try {
+            final Future<?> future = executor.submit(refresh);
+            assertTrue(preparing.await(10, TimeUnit.SECONDS));
+            for (int i = 0; i < 100; i++) {
+                readOld.run();
+            }
+            assertFalse(future.isDone());
+            publish.countDown();
+            future.get(10, TimeUnit.SECONDS);
+        } finally {
+            publish.countDown();
+            executor.shutdownNow();
+            assertTrue(executor.awaitTermination(10, TimeUnit.SECONDS));
+        }
+    }
+
+    private PluginData plugin(final String name, final int sort) {
+        return 
PluginData.builder().name(name).sort(sort).enabled(true).build();
+    }
+
+    private SelectorData selector(final String id, final int sort) {
+        return 
SelectorData.builder().id(id).pluginName("divide").sort(sort).build();
+    }
+
+    private RuleData rule(final String id, final int sort) {
+        return 
RuleData.builder().id(id).pluginName("divide").selectorId("selector").sort(sort).build();
+    }
 }
diff --git 
a/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriberTest.java
 
b/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriberTest.java
index 3df6109763..aa191d558c 100644
--- 
a/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriberTest.java
+++ 
b/shenyu-plugin/shenyu-plugin-base/src/test/java/org/apache/shenyu/plugin/base/cache/CommonPluginDataSubscriberTest.java
@@ -23,8 +23,10 @@ import 
org.apache.shenyu.common.config.ShenyuConfig.SelectorMatchCache;
 import org.apache.shenyu.common.dto.PluginData;
 import org.apache.shenyu.common.dto.RuleData;
 import org.apache.shenyu.common.dto.SelectorData;
+import org.apache.shenyu.common.enums.PluginHandlerEventEnum;
 import org.apache.shenyu.plugin.api.utils.SpringBeanUtils;
 import org.apache.shenyu.plugin.base.handler.PluginDataHandler;
+import org.junit.jupiter.api.AfterEach;
 import org.junit.jupiter.api.BeforeEach;
 import org.junit.jupiter.api.Test;
 import org.junit.jupiter.api.extension.ExtendWith;
@@ -40,8 +42,17 @@ import java.util.List;
 
 import static org.junit.jupiter.api.Assertions.assertEquals;
 import static org.junit.jupiter.api.Assertions.assertNotNull;
+import static org.junit.jupiter.api.Assertions.assertNotSame;
 import static org.junit.jupiter.api.Assertions.assertNull;
+import static org.junit.jupiter.api.Assertions.assertSame;
+import static org.junit.jupiter.api.Assertions.assertThrows;
+import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.Mockito.doAnswer;
+import static org.mockito.Mockito.doThrow;
 import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.Mockito.when;
 
 /**
  * Test cases for CommonPluginDataSubscriber.
@@ -69,12 +80,18 @@ public final class CommonPluginDataSubscriberTest {
     
     private BaseDataCache baseDataCache;
 
+    @Mock
+    private PluginDataHandler handler;
+
     @BeforeEach
     public void setup() {
         this.mockShenyuTrieConfig();
         ArrayList<PluginDataHandler> pluginDataHandlerList = 
Lists.newArrayList();
         commonPluginDataSubscriber = new 
CommonPluginDataSubscriber(pluginDataHandlerList, eventPublisher, new 
SelectorMatchCache(), new RuleMatchCache());
         baseDataCache = BaseDataCache.getInstance();
+        clearCaches();
+        when(handler.pluginNamed()).thenReturn("divide");
+        
commonPluginDataSubscriber.putExtendPluginDataHandler(List.of(handler));
     }
 
     @Test
@@ -261,4 +278,163 @@ public final class CommonPluginDataSubscriberTest {
         ConfigurableApplicationContext context = 
mock(ConfigurableApplicationContext.class);
         SpringBeanUtils.getInstance().setApplicationContext(context);
     }
+
+    @AfterEach
+    public void clearCaches() {
+        baseDataCache.cleanPluginData();
+        baseDataCache.cleanSelectorData();
+        baseDataCache.cleanRuleData();
+        MatchDataCache.getInstance().cleanSelectorData();
+        MatchDataCache.getInstance().cleanRuleDataData();
+    }
+
+    @Test
+    public void pluginRefreshPublishesTheCompleteBatch() {
+        final PluginData a = plugin("jwt", 1);
+        final PluginData old = plugin("divide", 2);
+        final PluginData updated = plugin("divide", 3);
+        final PluginData added = plugin("rewrite", 4);
+        baseDataCache.cachePluginData(a);
+        baseDataCache.cachePluginData(old);
+        final var previous = baseDataCache.getPluginMap();
+        doAnswer(invocation -> {
+            assertSame(previous, baseDataCache.getPluginMap());
+            assertSame(old, baseDataCache.obtainPluginData("divide"));
+            assertNull(baseDataCache.obtainPluginData("rewrite"));
+            verifyNoInteractions(eventPublisher);
+            return null;
+        }).when(handler).handlerPlugin(updated);
+        doAnswer(invocation -> {
+            assertSame(updated, baseDataCache.obtainPluginData("divide"));
+            assertSame(added, baseDataCache.obtainPluginData("rewrite"));
+            return null;
+        }).when(eventPublisher).publishEvent(any(PluginHandlerEvent.class));
+        commonPluginDataSubscriber.onPluginRefresh(List.of(updated, added));
+        assertNotSame(previous, baseDataCache.getPluginMap());
+        assertEquals(3, baseDataCache.getPluginMap().size());
+        assertSame(a, baseDataCache.obtainPluginData("jwt"));
+        assertSame(old, previous.get("divide"));
+        
verify(eventPublisher).publishEvent(org.mockito.ArgumentMatchers.argThat((PluginHandlerEvent
 event) ->
+                event.getSource() == updated && event.getPluginStateEnums() == 
PluginHandlerEventEnum.ENABLED));
+        
verify(eventPublisher).publishEvent(org.mockito.ArgumentMatchers.argThat((PluginHandlerEvent
 event) ->
+                event.getSource() == updated && event.getPluginStateEnums() == 
PluginHandlerEventEnum.SORTED));
+    }
+
+    @Test
+    public void selectorHandlerSeesThePublishedBatch() {
+        final SelectorData first = selector("a", 1);
+        final SelectorData second = selector("b", 2);
+        doAnswer(invocation -> {
+            assertEquals(List.of(first, second), 
baseDataCache.obtainSelectorData("divide"));
+            return null;
+        }).when(handler).handlerSelector(first);
+        commonPluginDataSubscriber.onSelectorRefresh(List.of(first, second));
+        verify(handler).handlerSelector(second);
+    }
+
+    @Test
+    public void ruleHandlerSeesThePublishedBatch() {
+        final RuleData first = rule("a", 1);
+        final RuleData second = rule("b", 2);
+        doAnswer(invocation -> {
+            assertEquals(List.of(first, second), 
baseDataCache.obtainRuleData("selector"));
+            return null;
+        }).when(handler).handlerRule(first);
+        commonPluginDataSubscriber.onRuleRefresh(List.of(first, second));
+        verify(handler).handlerRule(second);
+    }
+
+    @Test
+    public void emptyBatchesRetainAllMaps() {
+        baseDataCache.cachePluginData(plugin("divide", 1));
+        baseDataCache.cacheSelectData(selector("a", 1));
+        baseDataCache.cacheRuleData(rule("a", 1));
+        final var plugins = baseDataCache.getPluginMap();
+        final var selectors = baseDataCache.getSelectorMap();
+        final var rules = baseDataCache.getRuleMap();
+        commonPluginDataSubscriber.onPluginRefresh(List.of());
+        commonPluginDataSubscriber.onSelectorRefresh(List.of());
+        commonPluginDataSubscriber.onRuleRefresh(List.of());
+        assertSame(plugins, baseDataCache.getPluginMap());
+        assertSame(selectors, baseDataCache.getSelectorMap());
+        assertSame(rules, baseDataCache.getRuleMap());
+        verifyNoInteractions(eventPublisher);
+    }
+
+    @Test
+    public void selectorHandlerFailureDoesNotRollBackPublishedBatch() {
+        final SelectorData old = selector("b", 1);
+        final SelectorData updated = selector("b", 2);
+        final SelectorData added = selector("c", 3);
+        baseDataCache.cacheSelectData(old);
+        doThrow(new IllegalStateException("handler 
failed")).when(handler).handlerSelector(updated);
+        assertThrows(IllegalStateException.class, () -> 
commonPluginDataSubscriber.onSelectorRefresh(List.of(updated, added)));
+        assertEquals(List.of(updated, added), 
baseDataCache.obtainSelectorData("divide"));
+    }
+
+    @Test
+    public void ruleHandlerFailureDoesNotRollBackPublishedBatch() {
+        final RuleData old = rule("b", 1);
+        final RuleData updated = rule("b", 2);
+        final RuleData added = rule("c", 3);
+        baseDataCache.cacheRuleData(old);
+        doThrow(new IllegalStateException("handler 
failed")).when(handler).handlerRule(updated);
+        assertThrows(IllegalStateException.class, () -> 
commonPluginDataSubscriber.onRuleRefresh(List.of(updated, added)));
+        assertEquals(List.of(updated, added), 
baseDataCache.obtainRuleData("selector"));
+    }
+
+    @Test
+    public void refreshInvalidatesMatchingAndNegativeCacheEntries() {
+        final SelectorMatchCache selectorConfig = new SelectorMatchCache();
+        final RuleMatchCache ruleConfig = new RuleMatchCache();
+        selectorConfig.getCache().setEnabled(true);
+        ruleConfig.getCache().setEnabled(true);
+        commonPluginDataSubscriber = new 
CommonPluginDataSubscriber(List.of(handler), eventPublisher, selectorConfig, 
ruleConfig);
+        final MatchDataCache matches = MatchDataCache.getInstance();
+        matches.cacheSelectorData("/selector", selector("selector", 1), 100, 
100);
+        matches.cacheSelectorData("/empty", 
SelectorData.builder().pluginName("divide").build(), 100, 100);
+        matches.cacheRuleData("/selector", rule("a", 1), 100, 100);
+        matches.cacheRuleData("/empty", 
RuleData.builder().pluginName("divide").build(), 100, 100);
+        
commonPluginDataSubscriber.onSelectorRefresh(List.of(selector("selector", 2)));
+        assertNull(matches.obtainSelectorData("divide", "/selector"));
+        assertNull(matches.obtainSelectorData("divide", "/empty"));
+        assertNull(matches.obtainRuleData("divide", "/selector"));
+        assertNull(matches.obtainRuleData("divide", "/empty"));
+
+        matches.cacheRuleData("/rule", rule("a", 1), 100, 100);
+        matches.cacheRuleData("/empty", 
RuleData.builder().pluginName("divide").build(), 100, 100);
+        commonPluginDataSubscriber.onRuleRefresh(List.of(rule("a", 2)));
+        assertNull(matches.obtainRuleData("divide", "/rule"));
+        assertNull(matches.obtainRuleData("divide", "/empty"));
+
+        matches.cacheSelectorData("/selector", selector("a", 1), 100, 100);
+        matches.cacheRuleData("/rule", rule("a", 1), 100, 100);
+        commonPluginDataSubscriber.onPluginRefresh(List.of(plugin("divide", 
1)));
+        assertNull(matches.obtainSelectorData("divide", "/selector"));
+        assertNull(matches.obtainRuleData("divide", "/rule"));
+    }
+
+    @Test
+    public void disabledPluginEventIsPublishedAfterTheBatch() {
+        final PluginData disabled = 
PluginData.builder().name("divide").enabled(false).sort(1).build();
+        doAnswer(invocation -> {
+            assertSame(disabled, baseDataCache.obtainPluginData("divide"));
+            return null;
+        }).when(eventPublisher).publishEvent(any(PluginHandlerEvent.class));
+        commonPluginDataSubscriber.onPluginRefresh(List.of(disabled));
+        
verify(eventPublisher).publishEvent(org.mockito.ArgumentMatchers.argThat((PluginHandlerEvent
 event) ->
+                event.getSource() == disabled && event.getPluginStateEnums() 
== PluginHandlerEventEnum.DISABLED));
+    }
+
+    private PluginData plugin(final String name, final int sort) {
+        return 
PluginData.builder().name(name).sort(sort).enabled(true).build();
+    }
+
+    private SelectorData selector(final String id, final int sort) {
+        return 
SelectorData.builder().id(id).pluginName("divide").sort(sort).build();
+    }
+
+    private RuleData rule(final String id, final int sort) {
+        return 
RuleData.builder().id(id).pluginName("divide").selectorId("selector").sort(sort).build();
+    }
 }
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/PluginDataSubscriber.java
 
b/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/PluginDataSubscriber.java
index 321e98077f..ef52e2fd9a 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/PluginDataSubscriber.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-api/src/main/java/org/apache/shenyu/sync/data/api/PluginDataSubscriber.java
@@ -116,4 +116,34 @@ public interface PluginDataSubscriber {
      */
     default void refreshRuleDataSelf(List<RuleData> ruleDataList) {
     }
+
+    /**
+     * Refresh a batch, retaining the legacy behavior for custom subscribers.
+     *
+     * @param dataList the received data
+     */
+    default void onPluginRefresh(List<PluginData> dataList) {
+        refreshPluginDataSelf(dataList);
+        dataList.forEach(this::onSubscribe);
+    }
+
+    /**
+     * Refresh a batch, retaining the legacy behavior for custom subscribers.
+     *
+     * @param dataList the received data
+     */
+    default void onSelectorRefresh(List<SelectorData> dataList) {
+        refreshSelectorDataSelf(dataList);
+        dataList.forEach(this::onSelectorSubscribe);
+    }
+
+    /**
+     * Refresh a batch, retaining the legacy behavior for custom subscribers.
+     *
+     * @param dataList the received data
+     */
+    default void onRuleRefresh(List<RuleData> dataList) {
+        refreshRuleDataSelf(dataList);
+        dataList.forEach(this::onRuleSubscribe);
+    }
 }
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-api/src/test/java/org/apache/shenyu/sync/data/api/PluginDataSubscriberTest.java
 
b/shenyu-sync-data-center/shenyu-sync-data-api/src/test/java/org/apache/shenyu/sync/data/api/PluginDataSubscriberTest.java
new file mode 100644
index 0000000000..a9d3d78012
--- /dev/null
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-api/src/test/java/org/apache/shenyu/sync/data/api/PluginDataSubscriberTest.java
@@ -0,0 +1,81 @@
+/*
+ * 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.shenyu.sync.data.api;
+
+import org.apache.shenyu.common.dto.PluginData;
+import org.apache.shenyu.common.dto.RuleData;
+import org.apache.shenyu.common.dto.SelectorData;
+import org.junit.jupiter.api.Test;
+import org.mockito.InOrder;
+
+import java.util.List;
+
+import static org.mockito.Mockito.CALLS_REAL_METHODS;
+import static org.mockito.Mockito.inOrder;
+import static org.mockito.Mockito.mock;
+
+/**
+ * Existing subscribers keep their refresh and subscription callbacks.
+ */
+class PluginDataSubscriberTest {
+
+    @Test
+    void pluginRefreshRetainsLegacyOrder() {
+        PluginDataSubscriber subscriber = mock(PluginDataSubscriber.class, 
CALLS_REAL_METHODS);
+        PluginData first = PluginData.builder().name("first").build();
+        PluginData second = PluginData.builder().name("second").build();
+        List<PluginData> batch = List.of(first, second);
+        subscriber.onPluginRefresh(batch);
+        InOrder order = inOrder(subscriber);
+        order.verify(subscriber).onPluginRefresh(batch);
+        order.verify(subscriber).refreshPluginDataSelf(batch);
+        order.verify(subscriber).onSubscribe(first);
+        order.verify(subscriber).onSubscribe(second);
+        order.verifyNoMoreInteractions();
+    }
+
+    @Test
+    void selectorRefreshRetainsLegacyOrder() {
+        PluginDataSubscriber subscriber = mock(PluginDataSubscriber.class, 
CALLS_REAL_METHODS);
+        SelectorData first = SelectorData.builder().name("first").build();
+        SelectorData second = SelectorData.builder().name("second").build();
+        List<SelectorData> batch = List.of(first, second);
+        subscriber.onSelectorRefresh(batch);
+        InOrder order = inOrder(subscriber);
+        order.verify(subscriber).onSelectorRefresh(batch);
+        order.verify(subscriber).refreshSelectorDataSelf(batch);
+        order.verify(subscriber).onSelectorSubscribe(first);
+        order.verify(subscriber).onSelectorSubscribe(second);
+        order.verifyNoMoreInteractions();
+    }
+
+    @Test
+    void ruleRefreshRetainsLegacyOrder() {
+        PluginDataSubscriber subscriber = mock(PluginDataSubscriber.class, 
CALLS_REAL_METHODS);
+        RuleData first = RuleData.builder().name("first").build();
+        RuleData second = RuleData.builder().name("second").build();
+        List<RuleData> batch = List.of(first, second);
+        subscriber.onRuleRefresh(batch);
+        InOrder order = inOrder(subscriber);
+        order.verify(subscriber).onRuleRefresh(batch);
+        order.verify(subscriber).refreshRuleDataSelf(batch);
+        order.verify(subscriber).onRuleSubscribe(first);
+        order.verify(subscriber).onRuleSubscribe(second);
+        order.verifyNoMoreInteractions();
+    }
+}
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/PluginDataHandler.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/PluginDataHandler.java
index b55d7491a4..40e4b22719 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/PluginDataHandler.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/PluginDataHandler.java
@@ -41,8 +41,7 @@ public class PluginDataHandler extends 
AbstractDataHandler<PluginData> {
 
     @Override
     protected void doRefresh(final List<PluginData> dataList) {
-        pluginDataSubscriber.refreshPluginDataSelf(dataList);
-        dataList.forEach(pluginDataSubscriber::onSubscribe);
+        pluginDataSubscriber.onPluginRefresh(dataList);
     }
 
     @Override
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/RuleDataHandler.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/RuleDataHandler.java
index d9a47b5cc3..31e8ff697b 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/RuleDataHandler.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/RuleDataHandler.java
@@ -40,8 +40,7 @@ public class RuleDataHandler extends 
AbstractDataHandler<RuleData> {
 
     @Override
     protected void doRefresh(final List<RuleData> dataList) {
-        pluginDataSubscriber.refreshRuleDataSelf(dataList);
-        dataList.forEach(pluginDataSubscriber::onRuleSubscribe);
+        pluginDataSubscriber.onRuleRefresh(dataList);
     }
 
     @Override
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/SelectorDataHandler.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/SelectorDataHandler.java
index d33eb76344..8682485636 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/SelectorDataHandler.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/main/java/org/apache/shenyu/plugin/sync/data/websocket/handler/SelectorDataHandler.java
@@ -41,8 +41,7 @@ public class SelectorDataHandler extends 
AbstractDataHandler<SelectorData> {
 
     @Override
     protected void doRefresh(final List<SelectorData> dataList) {
-        pluginDataSubscriber.refreshSelectorDataSelf(dataList);
-        dataList.forEach(pluginDataSubscriber::onSelectorSubscribe);
+        pluginDataSubscriber.onSelectorRefresh(dataList);
     }
 
     @Override
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/client/ShenyuWebsocketClientTest.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/client/ShenyuWebsocketClientTest.java
index 82abc3b8e7..ad909e7a4a 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/client/ShenyuWebsocketClientTest.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/client/ShenyuWebsocketClientTest.java
@@ -114,20 +114,20 @@ public class ShenyuWebsocketClientTest {
 
     @Test
     public void testOnMessage() {
-        doNothing().when(pluginDataSubscriber).onSubscribe(any());
+        doNothing().when(pluginDataSubscriber).onPluginRefresh(any());
         String json = GsonUtils.getInstance().toJson(websocketData);
         shenyuWebsocketClient.onMessage(json);
-        verify(pluginDataSubscriber).onSubscribe(any());
+        verify(pluginDataSubscriber).onPluginRefresh(any());
     }
 
     @Test
     public void testOnMessageShouldIgnoreMalformedJsonAndHandleNextMessage() {
         Assertions.assertDoesNotThrow(() -> 
shenyuWebsocketClient.onMessage("{invalid json"));
 
-        doNothing().when(pluginDataSubscriber).onSubscribe(any());
+        doNothing().when(pluginDataSubscriber).onPluginRefresh(any());
         String json = GsonUtils.getInstance().toJson(websocketData);
         shenyuWebsocketClient.onMessage(json);
-        verify(pluginDataSubscriber).onSubscribe(any());
+        verify(pluginDataSubscriber).onPluginRefresh(any());
     }
 
     @Test
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/PluginDataHandlerTest.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/PluginDataHandlerTest.java
index f824648d5b..978d49030f 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/PluginDataHandlerTest.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/PluginDataHandlerTest.java
@@ -21,6 +21,8 @@ import com.google.gson.Gson;
 import org.apache.shenyu.common.dto.PluginData;
 import org.apache.shenyu.sync.data.api.PluginDataSubscriber;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
 
 import java.util.LinkedList;
 import java.util.List;
@@ -29,6 +31,8 @@ import static org.hamcrest.MatcherAssert.assertThat;
 import static org.hamcrest.core.Is.is;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.Mockito.verifyNoMoreInteractions;
 
 public final class PluginDataHandlerTest {
 
@@ -56,8 +60,24 @@ public final class PluginDataHandlerTest {
     public void testDoRefresh() {
         List<PluginData> pluginDataList = createFakePluginDataObjects(3);
         pluginDataHandler.doRefresh(pluginDataList);
-        verify(subscriber).refreshPluginDataSelf(pluginDataList);
-        pluginDataList.forEach(verify(subscriber)::onSubscribe);
+        verify(subscriber).onPluginRefresh(pluginDataList);
+    }
+
+
+    @ParameterizedTest
+    @ValueSource(strings = {"REFRESH", "MYSELF"})
+    void testRefreshEventsUseOnlyTheBatchCallback(final String eventType) {
+        List<PluginData> batch = createFakePluginDataObjects(2);
+        pluginDataHandler.handle(new Gson().toJson(batch), eventType);
+        verify(subscriber).onPluginRefresh(batch);
+        verifyNoMoreInteractions(subscriber);
+    }
+
+    @ParameterizedTest
+    @ValueSource(strings = {"REFRESH", "MYSELF"})
+    void testEmptyRefreshDoesNotCallSubscriber(final String eventType) {
+        pluginDataHandler.handle("[]", eventType);
+        verifyNoInteractions(subscriber);
     }
 
     @Test
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/RuleDataHandlerTest.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/RuleDataHandlerTest.java
index cf28d5fc6b..648483b3e5 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/RuleDataHandlerTest.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/RuleDataHandlerTest.java
@@ -22,6 +22,8 @@ import org.apache.shenyu.common.dto.ConditionData;
 import org.apache.shenyu.common.dto.RuleData;
 import org.apache.shenyu.sync.data.api.PluginDataSubscriber;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
 
 import java.util.Collections;
 import java.util.LinkedList;
@@ -31,6 +33,8 @@ import static org.hamcrest.MatcherAssert.assertThat;
 import static org.hamcrest.core.Is.is;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.Mockito.verifyNoMoreInteractions;
 
 public final class RuleDataHandlerTest {
 
@@ -61,8 +65,24 @@ public final class RuleDataHandlerTest {
     public void testDoRefresh() {
         List<RuleData> ruleDataList = createFakeRuleDateObjects(3);
         ruleDataHandler.doRefresh(ruleDataList);
-        verify(subscriber).refreshRuleDataSelf(ruleDataList);
-        ruleDataList.forEach(verify(subscriber)::onRuleSubscribe);
+        verify(subscriber).onRuleRefresh(ruleDataList);
+    }
+
+
+    @ParameterizedTest
+    @ValueSource(strings = {"REFRESH", "MYSELF"})
+    void testRefreshEventsUseOnlyTheBatchCallback(final String eventType) {
+        List<RuleData> batch = createFakeRuleDateObjects(2);
+        ruleDataHandler.handle(new Gson().toJson(batch), eventType);
+        verify(subscriber).onRuleRefresh(batch);
+        verifyNoMoreInteractions(subscriber);
+    }
+
+    @ParameterizedTest
+    @ValueSource(strings = {"REFRESH", "MYSELF"})
+    void testEmptyRefreshDoesNotCallSubscriber(final String eventType) {
+        ruleDataHandler.handle("[]", eventType);
+        verifyNoInteractions(subscriber);
     }
 
     @Test
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/SelectorDataHandlerTest.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/SelectorDataHandlerTest.java
index 1a64c43cdd..e58e7037ee 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/SelectorDataHandlerTest.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/SelectorDataHandlerTest.java
@@ -22,6 +22,8 @@ import org.apache.shenyu.common.dto.ConditionData;
 import org.apache.shenyu.common.dto.SelectorData;
 import org.apache.shenyu.sync.data.api.PluginDataSubscriber;
 import org.junit.jupiter.api.Test;
+import org.junit.jupiter.params.ParameterizedTest;
+import org.junit.jupiter.params.provider.ValueSource;
 
 import java.util.Collections;
 import java.util.LinkedList;
@@ -31,6 +33,8 @@ import static org.hamcrest.MatcherAssert.assertThat;
 import static org.hamcrest.core.Is.is;
 import static org.mockito.Mockito.mock;
 import static org.mockito.Mockito.verify;
+import static org.mockito.Mockito.verifyNoInteractions;
+import static org.mockito.Mockito.verifyNoMoreInteractions;
 
 public final class SelectorDataHandlerTest {
 
@@ -61,8 +65,25 @@ public final class SelectorDataHandlerTest {
     public void testDoRefresh() {
         List<SelectorData> selectorDataList = createFakeSelectorDataObjects(3);
         selectorDataHandler.doRefresh(selectorDataList);
-        verify(subscriber).refreshSelectorDataSelf(selectorDataList);
-        selectorDataList.forEach(verify(subscriber)::onSelectorSubscribe);
+        verify(subscriber).onSelectorRefresh(selectorDataList);
+    }
+
+
+    @ParameterizedTest
+    @ValueSource(strings = {"REFRESH", "MYSELF"})
+    void testRefreshEventsUseOnlyTheBatchCallback(final String eventType) {
+        List<SelectorData> batch = createFakeSelectorDataObjects(2);
+        batch.forEach(data -> data.setContinued(true));
+        selectorDataHandler.handle(new Gson().toJson(batch), eventType);
+        verify(subscriber).onSelectorRefresh(batch);
+        verifyNoMoreInteractions(subscriber);
+    }
+
+    @ParameterizedTest
+    @ValueSource(strings = {"REFRESH", "MYSELF"})
+    void testEmptyRefreshDoesNotCallSubscriber(final String eventType) {
+        selectorDataHandler.handle("[]", eventType);
+        verifyNoInteractions(subscriber);
     }
 
     @Test
diff --git 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/WebsocketDataHandlerTest.java
 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/WebsocketDataHandlerTest.java
index 8a527233ce..4dca4ea295 100644
--- 
a/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/WebsocketDataHandlerTest.java
+++ 
b/shenyu-sync-data-center/shenyu-sync-data-websocket/src/test/java/org/apache/shenyu/plugin/sync/data/websocket/handler/WebsocketDataHandlerTest.java
@@ -65,7 +65,7 @@ public final class WebsocketDataHandlerTest {
         String json = getJson();
         websocketDataHandler.executor(ConfigGroupEnum.PLUGIN, json, 
DataEventTypeEnum.REFRESH.name());
         List<PluginData> pluginDataList = new 
PluginDataHandler(pluginDataSubscriber).convert(json);
-        
Mockito.verify(pluginDataSubscriber).refreshPluginDataSelf(pluginDataList);
+        Mockito.verify(pluginDataSubscriber).onPluginRefresh(pluginDataList);
     }
 
     @Test
@@ -73,7 +73,7 @@ public final class WebsocketDataHandlerTest {
         String json = getJson();
         websocketDataHandler.executor(ConfigGroupEnum.PLUGIN, json, 
DataEventTypeEnum.MYSELF.name());
         List<PluginData> pluginDataList = new 
PluginDataHandler(pluginDataSubscriber).convert(json);
-        
Mockito.verify(pluginDataSubscriber).refreshPluginDataSelf(pluginDataList);
+        Mockito.verify(pluginDataSubscriber).onPluginRefresh(pluginDataList);
     }
 
     @Test

Reply via email to