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..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,6 +37,7 @@ 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 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 ConcurrentMap); + Assertions.assertTrue(cacheDataField.get(consulSyncDataService) instanceof ConcurrentMap); + } }