From c1303f68b6d5199fb28db025ae6a091eeca4fdc1 Mon Sep 17 00:00:00 2001 From: dengliming Date: Sun, 20 Sep 2026 21:27:05 +0800 Subject: [PATCH 1/2] fix(consul): use concurrent watcher state maps --- .../sync/data/consul/ConsulSyncDataService.java | 6 +++--- .../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 d45b3990ee2f..8b9d1d70580a 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.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 consulIndexes = new HashMap<>(); + private final Map consulIndexes = new ConcurrentHashMap<>(); - private final Map> cacheConsulDataKeyMap = new HashMap<>(); + private final Map> 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 9f6a5f9c5019..c4e69a0e1401 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.util.ArrayList; import java.util.List; import java.util.Map; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ScheduledThreadPoolExecutor; import java.util.function.BiConsumer; import java.util.function.Consumer; @@ -131,8 +132,19 @@ public void testWatchConfigKeyValues() throws NoSuchMethodException, IllegalAcce final Field consulIndexes = ConsulSyncDataService.class.getDeclaredField("consulIndexes"); consulIndexes.setAccessible(true); final Map consulIndexesSource = (Map) 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 ConcurrentHashMap); + Assertions.assertTrue(cacheDataField.get(consulSyncDataService) instanceof ConcurrentHashMap); + } } From d3f6fac3de9d1ed543a5fc5a2b8e5627891a01e7 Mon Sep 17 00:00:00 2001 From: dengliming Date: Mon, 21 Sep 2026 18:13:11 +0800 Subject: [PATCH 2/2] test(consul): assert concurrent map contract --- .../shenyu/sync/data/consul/ConsulSyncDataServiceTest.java | 6 +++--- 1 file changed, 3 insertions(+), 3 deletions(-) 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 c4e69a0e1401..606e78df3956 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,7 +37,7 @@ import java.util.ArrayList; import java.util.List; import java.util.Map; -import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ScheduledThreadPoolExecutor; import java.util.function.BiConsumer; import java.util.function.Consumer; @@ -144,7 +144,7 @@ public void testWatcherStateUsesConcurrentMaps() throws NoSuchFieldException, Il Field cacheDataField = ConsulSyncDataService.class.getDeclaredField("cacheConsulDataKeyMap"); cacheDataField.setAccessible(true); - Assertions.assertTrue(consulIndexesField.get(consulSyncDataService) instanceof ConcurrentHashMap); - Assertions.assertTrue(cacheDataField.get(consulSyncDataService) instanceof ConcurrentHashMap); + Assertions.assertTrue(consulIndexesField.get(consulSyncDataService) instanceof ConcurrentMap); + Assertions.assertTrue(cacheDataField.get(consulSyncDataService) instanceof ConcurrentMap); } }