From e9587aadde33d3084784aeb0ae74d7f45b510b23 Mon Sep 17 00:00:00 2001 From: liuhy Date: Sun, 9 Aug 2026 20:05:03 -0700 Subject: [PATCH] [ISSUE #10892] Tolerate concurrent topic cache eviction Signed-off-by: liuhy --- .../impl/producer/DefaultMQProducerImpl.java | 6 ++- .../selector/DefaultMQProducerImplTest.java | 31 ++++++++++++++ .../service/route/TopicRouteService.java | 6 ++- .../route/ClusterTopicRouteServiceTest.java | 41 +++++++++++++++++++ 4 files changed, 80 insertions(+), 4 deletions(-) diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java b/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java index 9ad5fcef4dc..b06b0873f1c 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/producer/DefaultMQProducerImpl.java @@ -21,6 +21,7 @@ import java.util.ArrayList; import java.util.Arrays; import java.util.HashSet; +import java.util.Iterator; import java.util.List; import java.util.Random; import java.util.Set; @@ -178,10 +179,11 @@ public String resolve(String name) { }, serviceDetector); } private Optional pickTopic() { - if (topicPublishInfoTable.isEmpty()) { + Iterator iterator = topicPublishInfoTable.keySet().iterator(); + if (!iterator.hasNext()) { return Optional.empty(); } - return Optional.of(topicPublishInfoTable.keySet().iterator().next()); + return Optional.of(iterator.next()); } public void registerCheckForbiddenHook(CheckForbiddenHook checkForbiddenHook) { this.checkForbiddenHookList.add(checkForbiddenHook); diff --git a/client/src/test/java/org/apache/rocketmq/client/producer/selector/DefaultMQProducerImplTest.java b/client/src/test/java/org/apache/rocketmq/client/producer/selector/DefaultMQProducerImplTest.java index 77a83af19c0..cf60e017e54 100644 --- a/client/src/test/java/org/apache/rocketmq/client/producer/selector/DefaultMQProducerImplTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/producer/selector/DefaultMQProducerImplTest.java @@ -54,6 +54,7 @@ import java.util.ArrayList; import java.util.Collections; import java.util.List; +import java.util.Optional; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; import java.util.concurrent.ExecutorService; @@ -222,6 +223,22 @@ public void testInitTopicRoute() throws NoSuchMethodException, InvocationTargetE method.invoke(defaultMQProducerImpl); } + @Test + public void testPickTopicToleratesConcurrentCacheEviction() throws Exception { + ClearingConcurrentMap topicPublishInfoTable = new ClearingConcurrentMap<>(); + topicPublishInfoTable.put(defaultTopic, mock(TopicPublishInfo.class)); + setField(defaultMQProducerImpl, "topicPublishInfoTable", topicPublishInfoTable); + + Method method = DefaultMQProducerImpl.class.getDeclaredMethod("pickTopic"); + method.setAccessible(true); + + Optional topic = (Optional) method.invoke(defaultMQProducerImpl); + + // ConcurrentHashMap iterators are weakly consistent, so a candidate observed before + // eviction is valid. The important contract is that selection never throws. + assertNotNull(topic); + } + @Test public void assertFetchPublishMessageQueues() throws MQClientException { List actual = defaultMQProducerImpl.fetchPublishMessageQueues(defaultTopic); @@ -382,4 +399,18 @@ private void setField(final Object target, final String fieldName, final Object field.setAccessible(true); field.set(target, newValue); } + + private static class ClearingConcurrentMap extends ConcurrentHashMap { + private boolean clearOnFirstIsEmpty = true; + + @Override + public boolean isEmpty() { + boolean empty = super.isEmpty(); + if (clearOnFirstIsEmpty) { + clearOnFirstIsEmpty = false; + super.clear(); + } + return empty; + } + } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/TopicRouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/TopicRouteService.java index dae30057461..128d4215a4f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/TopicRouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/TopicRouteService.java @@ -22,6 +22,7 @@ import com.google.common.annotations.VisibleForTesting; import java.time.Duration; import java.util.ArrayList; +import java.util.Iterator; import java.util.List; import java.util.Optional; import java.util.concurrent.ThreadPoolExecutor; @@ -134,10 +135,11 @@ public String resolve(String name) { // pickup one topic in the topic cache private Optional pickTopic() { - if (topicCache.asMap().isEmpty()) { + Iterator iterator = topicCache.asMap().keySet().iterator(); + if (!iterator.hasNext()) { return Optional.empty(); } - return Optional.of(topicCache.asMap().keySet().iterator().next()); + return Optional.of(iterator.next()); } protected void init() { diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/service/route/ClusterTopicRouteServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/service/route/ClusterTopicRouteServiceTest.java index 15d83483b9d..7eee933681f 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/service/route/ClusterTopicRouteServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/service/route/ClusterTopicRouteServiceTest.java @@ -22,8 +22,12 @@ import com.github.benmanes.caffeine.cache.LoadingCache; import com.google.common.net.HostAndPort; +import java.lang.reflect.Field; +import java.lang.reflect.Method; import java.util.HashMap; import java.util.List; +import java.util.Optional; +import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ThreadPoolExecutor; import java.util.concurrent.TimeUnit; import java.util.concurrent.atomic.AtomicBoolean; @@ -50,6 +54,7 @@ import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.eq; +import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; public class ClusterTopicRouteServiceTest extends BaseServiceTest { @@ -126,6 +131,22 @@ public void testGetTopicRouteForProxy() throws Throwable { assertEquals(addressList, proxyTopicRouteData.getBrokerDatas().get(0).getBrokerAddrs().get(MixAll.MASTER_ID)); } + @Test + public void testPickTopicToleratesConcurrentCacheEviction() throws Exception { + LoadingCache topicCache = mock(LoadingCache.class); + ClearingConcurrentMap topicCacheMap = new ClearingConcurrentMap<>(); + topicCacheMap.put(TOPIC, mock(MessageQueueView.class)); + when(topicCache.asMap()).thenReturn(topicCacheMap); + setField(topicRouteService, "topicCache", topicCache); + + Method method = TopicRouteService.class.getDeclaredMethod("pickTopic"); + method.setAccessible(true); + + Optional topic = (Optional) method.invoke(topicRouteService); + + assertNotNull(topic); + } + @Test public void testTopicRouteCaffeineCache() throws InterruptedException { String key = "abc"; @@ -164,4 +185,24 @@ public void testTopicRouteCaffeineCache() throws InterruptedException { TimeUnit.SECONDS.sleep(5); assertThat(value).isEqualTo(topicCache.get(key)); } + + private static void setField(Object target, String fieldName, Object value) throws Exception { + Field field = TopicRouteService.class.getDeclaredField(fieldName); + field.setAccessible(true); + field.set(target, value); + } + + private static class ClearingConcurrentMap extends ConcurrentHashMap { + private boolean clearOnFirstIsEmpty = true; + + @Override + public boolean isEmpty() { + boolean empty = super.isEmpty(); + if (clearOnFirstIsEmpty) { + clearOnFirstIsEmpty = false; + super.clear(); + } + return empty; + } + } }