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 5dbc24d4dd fix(consul): use concurrent watcher state maps (#7145)
5dbc24d4dd is described below
commit 5dbc24d4ddf6c774fdf6008563763009ee901545
Author: Liming Deng <[email protected]>
AuthorDate: Tue Sep 22 10:51:32 2026 +0800
fix(consul): use concurrent watcher state maps (#7145)
* fix(consul): use concurrent watcher state maps
* test(consul): assert concurrent map contract
---------
Co-authored-by: shown <[email protected]>
Co-authored-by: aias00 <[email protected]>
---
.../shenyu/sync/data/consul/ConsulSyncDataService.java | 6 +++---
.../shenyu/sync/data/consul/ConsulSyncDataServiceTest.java | 14 +++++++++++++-
2 files changed, 16 insertions(+), 4 deletions(-)
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-consul/src/main/java/org/apache/shenyu/sync/data/consul/ConsulSyncDataService.java
b/shenyu-sync-data-center/shenyu-sync-data-consul/src/main/java/org/apache/shenyu/sync/data/consul/ConsulSyncDataService.java
index d45b3990ee..8b9d1d7058 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-consul/src/main/java/org/apache/shenyu/sync/data/consul/ConsulSyncDataService.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-consul/src/main/java/org/apache/shenyu/sync/data/consul/ConsulSyncDataService.java
@@ -39,10 +39,10 @@ import
org.apache.shenyu.sync.data.core.AbstractPathDataSyncService;
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
-import java.util.HashMap;
import java.util.List;
import java.util.Map;
import java.util.Objects;
+import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.ScheduledThreadPoolExecutor;
import java.util.concurrent.TimeUnit;
import java.util.function.BiConsumer;
@@ -60,9 +60,9 @@ public class ConsulSyncDataService extends
AbstractPathDataSyncService {
*/
private static final Logger LOG =
LoggerFactory.getLogger(ConsulSyncDataService.class);
- private final Map<String, Long> consulIndexes = new HashMap<>();
+ private final Map<String, Long> consulIndexes = new ConcurrentHashMap<>();
- private final Map<String, List<ConsulData>> cacheConsulDataKeyMap = new
HashMap<>();
+ private final Map<String, List<ConsulData>> cacheConsulDataKeyMap = new
ConcurrentHashMap<>();
private final ScheduledThreadPoolExecutor executor;
diff --git
a/shenyu-sync-data-center/shenyu-sync-data-consul/src/test/java/org/apache/shenyu/sync/data/consul/ConsulSyncDataServiceTest.java
b/shenyu-sync-data-center/shenyu-sync-data-consul/src/test/java/org/apache/shenyu/sync/data/consul/ConsulSyncDataServiceTest.java
index 9f6a5f9c50..606e78df39 100644
---
a/shenyu-sync-data-center/shenyu-sync-data-consul/src/test/java/org/apache/shenyu/sync/data/consul/ConsulSyncDataServiceTest.java
+++
b/shenyu-sync-data-center/shenyu-sync-data-consul/src/test/java/org/apache/shenyu/sync/data/consul/ConsulSyncDataServiceTest.java
@@ -37,6 +37,7 @@ import java.lang.reflect.Method;
import java.util.ArrayList;
import java.util.List;
import java.util.Map;
+import java.util.concurrent.ConcurrentMap;
import java.util.concurrent.ScheduledThreadPoolExecutor;
import java.util.function.BiConsumer;
import java.util.function.Consumer;
@@ -131,8 +132,19 @@ public final class ConsulSyncDataServiceTest {
final Field consulIndexes =
ConsulSyncDataService.class.getDeclaredField("consulIndexes");
consulIndexes.setAccessible(true);
final Map<String, Long> consulIndexesSource = (Map<String, Long>)
consulIndexes.get(consulSyncDataService);
- consulIndexesSource.put("/null", null);
+ consulIndexesSource.remove(watchPathRoot);
when(response.getConsulIndex()).thenReturn(2L);
Assertions.assertDoesNotThrow(() ->
watchConfigKeyValues.invoke(consulSyncDataService, watchPathRoot,
updateHandler, deleteHandler));
}
+
+ @Test
+ public void testWatcherStateUsesConcurrentMaps() throws
NoSuchFieldException, IllegalAccessException {
+ Field consulIndexesField =
ConsulSyncDataService.class.getDeclaredField("consulIndexes");
+ consulIndexesField.setAccessible(true);
+ Field cacheDataField =
ConsulSyncDataService.class.getDeclaredField("cacheConsulDataKeyMap");
+ cacheDataField.setAccessible(true);
+
+ Assertions.assertTrue(consulIndexesField.get(consulSyncDataService)
instanceof ConcurrentMap);
+ Assertions.assertTrue(cacheDataField.get(consulSyncDataService)
instanceof ConcurrentMap);
+ }
}