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 fcf0979574 fix: use thread-safe maps for registry instance watchers 
(#6989)
fcf0979574 is described below

commit fcf0979574a18f63a77c86ba75eceffbba8930a5
Author: hengyuss <[email protected]>
AuthorDate: Fri Sep 4 17:05:11 2026 +0800

    fix: use thread-safe maps for registry instance watchers (#6989)
    
    Co-authored-by: aias00 <[email protected]>
---
 .../registry/apollo/ApolloInstanceRegisterRepository.java | 12 +++++++-----
 .../registry/consul/ConsulInstanceRegisterRepository.java | 15 ++++++++-------
 .../registry/etcd/EtcdInstanceRegisterRepository.java     |  8 ++++----
 .../zookeeper/ZookeeperInstanceRegisterRepository.java    | 11 ++++++-----
 4 files changed, 25 insertions(+), 21 deletions(-)

diff --git 
a/shenyu-registry/shenyu-registry-apollo/src/main/java/org/apache/shenyu/registry/apollo/ApolloInstanceRegisterRepository.java
 
b/shenyu-registry/shenyu-registry-apollo/src/main/java/org/apache/shenyu/registry/apollo/ApolloInstanceRegisterRepository.java
index 6c5440a130..730bb6f0b0 100644
--- 
a/shenyu-registry/shenyu-registry-apollo/src/main/java/org/apache/shenyu/registry/apollo/ApolloInstanceRegisterRepository.java
+++ 
b/shenyu-registry/shenyu-registry-apollo/src/main/java/org/apache/shenyu/registry/apollo/ApolloInstanceRegisterRepository.java
@@ -33,12 +33,13 @@ import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
 import java.net.URI;
-import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.Objects;
 import java.util.Optional;
 import java.util.Properties;
 import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
 import java.util.function.Function;
 import java.util.stream.Collectors;
 
@@ -62,7 +63,7 @@ public class ApolloInstanceRegisterRepository implements 
ShenyuInstanceRegisterR
 
     private final Map<String, ConfigChangeListener> configChangeListenerMap = 
Maps.newConcurrentMap();
 
-    private final Map<String, List<InstanceEntity>> watcherInstanceRegisterMap 
= new HashMap<>();
+    private final Map<String, List<InstanceEntity>> watcherInstanceRegisterMap 
= new ConcurrentHashMap<>();
 
     private String namespace;
 
@@ -125,10 +126,11 @@ public class ApolloInstanceRegisterRepository implements 
ShenyuInstanceRegisterR
                     instanceEntity.setUri(getURI(x, instanceEntity.getPort(), 
instanceEntity.getHost()));
                     return instanceEntity;
                 }).collect(Collectors.toList());
-        Map<String, String> childrenList = new HashMap<>();
+        Map<String, String> childrenList = new ConcurrentHashMap<>();
 
-        if (watcherInstanceRegisterMap.containsKey(selectKey)) {
-            return watcherInstanceRegisterMap.get(selectKey);
+        final List<InstanceEntity> cachedInstances = 
watcherInstanceRegisterMap.get(selectKey);
+        if (Objects.nonNull(cachedInstances)) {
+            return cachedInstances;
         }
 
         configService.getPropertyNames().forEach(key -> {
diff --git 
a/shenyu-registry/shenyu-registry-consul/src/main/java/org/apache/shenyu/registry/consul/ConsulInstanceRegisterRepository.java
 
b/shenyu-registry/shenyu-registry-consul/src/main/java/org/apache/shenyu/registry/consul/ConsulInstanceRegisterRepository.java
index b8f56fc231..73ef786628 100644
--- 
a/shenyu-registry/shenyu-registry-consul/src/main/java/org/apache/shenyu/registry/consul/ConsulInstanceRegisterRepository.java
+++ 
b/shenyu-registry/shenyu-registry-consul/src/main/java/org/apache/shenyu/registry/consul/ConsulInstanceRegisterRepository.java
@@ -41,13 +41,13 @@ import java.net.URI;
 import java.util.ArrayList;
 import java.util.Arrays;
 import java.util.Collections;
-import java.util.HashMap;
-import java.util.HashSet;
 import java.util.List;
 import java.util.Map;
+import java.util.Objects;
 import java.util.Optional;
 import java.util.Properties;
 import java.util.Set;
+import java.util.concurrent.ConcurrentHashMap;
 import java.util.concurrent.ScheduledFuture;
 import java.util.concurrent.ScheduledThreadPoolExecutor;
 import java.util.concurrent.TimeUnit;
@@ -73,7 +73,7 @@ public class ConsulInstanceRegisterRepository implements 
ShenyuInstanceRegisterR
 
     private final AtomicBoolean running = new AtomicBoolean(false);
 
-    private final Map<String, Long> consulIndexes = new HashMap<>();
+    private final Map<String, Long> consulIndexes = new ConcurrentHashMap<>();
 
     private String token;
 
@@ -87,9 +87,9 @@ public class ConsulInstanceRegisterRepository implements 
ShenyuInstanceRegisterR
 
     private TtlScheduler ttlScheduler;
 
-    private final Map<String, List<InstanceEntity>> watcherInstanceRegisterMap 
= new HashMap<>();
+    private final Map<String, List<InstanceEntity>> watcherInstanceRegisterMap 
= new ConcurrentHashMap<>();
 
-    private final Set<String> watchSelectKeySet = new HashSet<>();
+    private final Set<String> watchSelectKeySet = 
ConcurrentHashMap.newKeySet();
 
     @Override
     public void init(final RegisterConfig config) {
@@ -157,8 +157,9 @@ public class ConsulInstanceRegisterRepository implements 
ShenyuInstanceRegisterR
 
     @Override
     public List<InstanceEntity> selectInstances(final String selectKey) {
-        if (watcherInstanceRegisterMap.containsKey(selectKey)) {
-            return watcherInstanceRegisterMap.get(selectKey);
+        final List<InstanceEntity> cachedInstances = 
watcherInstanceRegisterMap.get(selectKey);
+        if (Objects.nonNull(cachedInstances)) {
+            return cachedInstances;
         }
         this.watcherStart(selectKey);
         final List<InstanceEntity> healthServices = 
this.getHealthServices(selectKey, "-1");
diff --git 
a/shenyu-registry/shenyu-registry-etcd/src/main/java/org/apache/shenyu/registry/etcd/EtcdInstanceRegisterRepository.java
 
b/shenyu-registry/shenyu-registry-etcd/src/main/java/org/apache/shenyu/registry/etcd/EtcdInstanceRegisterRepository.java
index 111a6528f6..0acb3f73de 100644
--- 
a/shenyu-registry/shenyu-registry-etcd/src/main/java/org/apache/shenyu/registry/etcd/EtcdInstanceRegisterRepository.java
+++ 
b/shenyu-registry/shenyu-registry-etcd/src/main/java/org/apache/shenyu/registry/etcd/EtcdInstanceRegisterRepository.java
@@ -37,11 +37,11 @@ import org.slf4j.LoggerFactory;
 
 import java.net.URI;
 import java.nio.charset.StandardCharsets;
-import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
 import java.util.Properties;
+import java.util.concurrent.ConcurrentHashMap;
 import java.util.function.Function;
 import java.util.stream.Collectors;
 
@@ -57,7 +57,7 @@ public class EtcdInstanceRegisterRepository implements 
ShenyuInstanceRegisterRep
 
     private EtcdClient client;
 
-    private final Map<String, List<InstanceEntity>> watcherInstanceRegisterMap 
= new HashMap<>();
+    private final Map<String, List<InstanceEntity>> watcherInstanceRegisterMap 
= new ConcurrentHashMap<>();
 
     private final Multimap<String, Watch.Watcher> watchCache = 
ArrayListMultimap.create();
 
@@ -95,10 +95,10 @@ public class EtcdInstanceRegisterRepository implements 
ShenyuInstanceRegisterRep
                     instanceEntity.setUri(getURI(x, instanceEntity.getPort(), 
instanceEntity.getHost()));
                     return instanceEntity;
                 }).collect(Collectors.toList());
-        if (watcherInstanceRegisterMap.containsKey(selectKey)) {
+        if (Objects.nonNull(watcherInstanceRegisterMap.get(selectKey))) {
             return 
getInstanceRegisterFun.apply(client.getKeysMapByPrefix(watchKey));
         }
-        Map<String, String> serverNodes = client.getKeysMapByPrefix(watchKey);
+        Map<String, String> serverNodes = new 
ConcurrentHashMap<>(client.getKeysMapByPrefix(watchKey));
         this.client.watchKeyChanges(watchKey, Watch.listener(response -> {
             for (WatchEvent event : response.getEvents()) {
                 String value = 
event.getKeyValue().getValue().toString(StandardCharsets.UTF_8);
diff --git 
a/shenyu-registry/shenyu-registry-zookeeper/src/main/java/org/apache/shenyu/registry/zookeeper/ZookeeperInstanceRegisterRepository.java
 
b/shenyu-registry/shenyu-registry-zookeeper/src/main/java/org/apache/shenyu/registry/zookeeper/ZookeeperInstanceRegisterRepository.java
index f548c211b1..4404c440b8 100644
--- 
a/shenyu-registry/shenyu-registry-zookeeper/src/main/java/org/apache/shenyu/registry/zookeeper/ZookeeperInstanceRegisterRepository.java
+++ 
b/shenyu-registry/shenyu-registry-zookeeper/src/main/java/org/apache/shenyu/registry/zookeeper/ZookeeperInstanceRegisterRepository.java
@@ -44,11 +44,11 @@ import org.slf4j.LoggerFactory;
 import java.net.URI;
 import java.nio.charset.StandardCharsets;
 import java.util.Collections;
-import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
 import java.util.Objects;
 import java.util.Properties;
+import java.util.concurrent.ConcurrentHashMap;
 import java.util.function.Function;
 import java.util.stream.Collectors;
 
@@ -64,11 +64,11 @@ public class ZookeeperInstanceRegisterRepository implements 
ShenyuInstanceRegist
 
     private String watchPath;
 
-    private final Map<String, String> nodeDataMap = new HashMap<>();
+    private final Map<String, String> nodeDataMap = new ConcurrentHashMap<>();
 
     private final Multimap<String, CuratorCache> cacheMap = 
ArrayListMultimap.create();
 
-    private final Map<String, List<InstanceEntity>> watcherInstanceRegisterMap 
= new HashMap<>();
+    private final Map<String, List<InstanceEntity>> watcherInstanceRegisterMap 
= new ConcurrentHashMap<>();
 
     @Override
     public void init(final RegisterConfig config) {
@@ -142,8 +142,9 @@ public class ZookeeperInstanceRegisterRepository implements 
ShenyuInstanceRegist
                 return instanceEntity;
             }).collect(Collectors.toList());
 
-            if (watcherInstanceRegisterMap.containsKey(selectKey)) {
-                return watcherInstanceRegisterMap.get(selectKey);
+            final List<InstanceEntity> cachedInstances = 
watcherInstanceRegisterMap.get(selectKey);
+            if (Objects.nonNull(cachedInstances)) {
+                return cachedInstances;
             }
 
             List<String> childrenPathList = 
client.subscribeChildrenChanges(watchKey, new CuratorWatcher() {

Reply via email to