This is an automated email from the ASF dual-hosted git repository.
yuluo-yx 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 34ec5f73f8 fix(sync): process zookeeper delete events (#7102)
34ec5f73f8 is described below
commit 34ec5f73f8bc7113dfdec6f61cf4761eea323cab
Author: Liming Deng <[email protected]>
AuthorDate: Thu Sep 24 10:03:06 2026 +0800
fix(sync): process zookeeper delete events (#7102)
Co-authored-by: aias00 <[email protected]>
Co-authored-by: shown <[email protected]>
---
.../data/zookeeper/ZookeeperSyncDataService.java | 73 +++++++++++++---------
.../zookeeper/ZookeeperSyncDataServiceTest.java | 25 ++++++++
2 files changed, 68 insertions(+), 30 deletions(-)
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/main/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataService.java
b/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/main/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataService.java
index 7bda998391..dab1b2c2aa 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/main/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataService.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/main/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataService.java
@@ -19,6 +19,8 @@ package org.apache.shenyu.sync.data.zookeeper;
import com.google.common.base.Strings;
import org.apache.commons.lang3.StringUtils;
+import org.apache.curator.framework.recipes.cache.ChildData;
+import org.apache.curator.framework.recipes.cache.CuratorCacheListener;
import org.apache.shenyu.common.config.ShenyuConfig;
import org.apache.shenyu.common.constant.Constants;
import org.apache.shenyu.common.constant.DefaultPathConstants;
@@ -82,37 +84,48 @@ public class ZookeeperSyncDataService extends
AbstractPathDataSyncService {
private void watcherData0(final String registerPath) {
String configNamespace = Constants.PATH_SEPARATOR +
shenyuConfig.getNamespace();
- zkClient.addCuratorCache(registerPath, (type, oldData, data) -> {
- if (Objects.isNull(data) || Objects.isNull(data.getData())) {
- return;
- }
- String path = data.getPath();
- if (Strings.isNullOrEmpty(path)) {
- return;
- }
- // if not uri register path, return.
- if (!path.contains(registerPath)) {
- return;
- }
- if (!StringUtils.containsIgnoreCase(path, configNamespace)) {
- return;
- }
+ zkClient.addCuratorCache(registerPath,
+ (type, oldData, data) -> handleEvent(type, oldData, data,
registerPath, configNamespace));
+ }
- EventType eventType = EventType.PUT;
- switch (type) {
- case NODE_DELETED:
- eventType = EventType.DELETE;
- break;
- case NODE_CREATED:
- case NODE_CHANGED:
- eventType = EventType.PUT;
- break;
- default:
- break;
- }
- final String updateData = new String(data.getData(),
StandardCharsets.UTF_8);
- this.event(configNamespace, path, updateData, registerPath,
eventType);
- });
+ private void handleEvent(final CuratorCacheListener.Type type, final
ChildData oldData, final ChildData data,
+ final String registerPath, final String
configNamespace) {
+ final ChildData eventData;
+ final EventType eventType;
+ final String updateData;
+ switch (type) {
+ case NODE_DELETED:
+ eventData = oldData;
+ eventType = EventType.DELETE;
+ updateData = null;
+ break;
+ case NODE_CREATED:
+ case NODE_CHANGED:
+ eventData = data;
+ eventType = EventType.PUT;
+ if (Objects.isNull(eventData) ||
Objects.isNull(eventData.getData())) {
+ return;
+ }
+ updateData = new String(eventData.getData(),
StandardCharsets.UTF_8);
+ break;
+ default:
+ return;
+ }
+ if (Objects.isNull(eventData)) {
+ return;
+ }
+ String path = eventData.getPath();
+ if (Strings.isNullOrEmpty(path)) {
+ return;
+ }
+ // if not uri register path, return.
+ if (!path.contains(registerPath)) {
+ return;
+ }
+ if (!StringUtils.containsIgnoreCase(path, configNamespace)) {
+ return;
+ }
+ this.event(configNamespace, path, updateData, registerPath, eventType);
}
@Override
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/test/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataServiceTest.java
b/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/test/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataServiceTest.java
index b572c99b1c..78d906d00f 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/test/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataServiceTest.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-zookeeper/src/test/java/org/apache/shenyu/sync/data/zookeeper/ZookeeperSyncDataServiceTest.java
@@ -19,6 +19,7 @@ package org.apache.shenyu.sync.data.zookeeper;
import org.apache.curator.framework.CuratorFramework;
import org.apache.curator.framework.recipes.cache.ChildData;
+import org.apache.curator.framework.recipes.cache.CuratorCacheListener;
import org.apache.curator.framework.recipes.cache.TreeCacheEvent;
import org.apache.curator.framework.recipes.cache.TreeCacheListener;
import org.apache.shenyu.common.config.ShenyuConfig;
@@ -30,6 +31,7 @@ import org.apache.shenyu.sync.data.api.PluginDataSubscriber;
import org.apache.shenyu.sync.data.api.ProxySelectorDataSubscriber;
import org.apache.shenyu.sync.data.api.DiscoveryUpstreamDataSubscriber;
import org.junit.jupiter.api.Test;
+import org.mockito.ArgumentCaptor;
import java.util.ArrayList;
import java.util.Collections;
@@ -37,12 +39,35 @@ import java.util.List;
import java.util.Objects;
import static org.mockito.ArgumentMatchers.any;
+import static org.mockito.ArgumentMatchers.anyString;
+import static org.mockito.ArgumentMatchers.argThat;
import static org.mockito.Mockito.doAnswer;
import static org.mockito.Mockito.mock;
+import static org.mockito.Mockito.times;
+import static org.mockito.Mockito.verify;
import static org.mockito.Mockito.when;
public final class ZookeeperSyncDataServiceTest {
+ @Test
+ public void testDeletedNodeUsesOldData() {
+ ZookeeperClient zkClient = mock(ZookeeperClient.class);
+ PluginDataSubscriber pluginDataSubscriber =
mock(PluginDataSubscriber.class);
+ ShenyuConfig shenyuConfig = mock(ShenyuConfig.class);
+
when(shenyuConfig.getNamespace()).thenReturn(Constants.SYS_DEFAULT_NAMESPACE_ID);
+ new ZookeeperSyncDataService(shenyuConfig, zkClient,
pluginDataSubscriber, Collections.emptyList(),
+ Collections.emptyList(), Collections.emptyList(),
Collections.emptyList());
+ ArgumentCaptor<CuratorCacheListener> listenerCaptor =
ArgumentCaptor.forClass(CuratorCacheListener.class);
+ verify(zkClient, times(7)).addCuratorCache(anyString(),
listenerCaptor.capture());
+
+ String pluginPath = Constants.PATH_SEPARATOR +
Constants.SYS_DEFAULT_NAMESPACE_ID + "/shenyu/plugin/divide";
+ ChildData oldData = new ChildData(pluginPath, null, "{}".getBytes());
+ listenerCaptor.getAllValues().forEach(listener ->
+ listener.event(CuratorCacheListener.Type.NODE_DELETED,
oldData, null));
+
+ verify(pluginDataSubscriber).unSubscribe(argThat(pluginData ->
"divide".equals(pluginData.getName())));
+ }
+
@Test
public void testZookeeperInstanceRegisterRepository() throws Exception {