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() {