From a65c345b348517feb1f8f10ae38e97a24d0e082e Mon Sep 17 00:00:00 2001 From: yuluo-yx Date: Sat, 8 Aug 2026 15:12:03 +0800 Subject: [PATCH] [ISSUE #10851] fix(broker): avoid cold threshold sort overflow --- .../broker/coldctr/ColdDataCgCtrService.java | 2 +- .../coldctr/ColdDataCgCtrServiceTest.java | 25 +++++++++++++++++++ 2 files changed, 26 insertions(+), 1 deletion(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/coldctr/ColdDataCgCtrService.java b/broker/src/main/java/org/apache/rocketmq/broker/coldctr/ColdDataCgCtrService.java index 5b8b2fb9cec..36af5ab574b 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/coldctr/ColdDataCgCtrService.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/coldctr/ColdDataCgCtrService.java @@ -145,7 +145,7 @@ private void sortAndDecelerate() { configMapList.sort(new Comparator>() { @Override public int compare(Entry o1, Entry o2) { - return (int)(o2.getValue() - o1.getValue()); + return Long.compare(o2.getValue(), o1.getValue()); } }); Iterator> iterator = configMapList.iterator(); diff --git a/broker/src/test/java/org/apache/rocketmq/broker/coldctr/ColdDataCgCtrServiceTest.java b/broker/src/test/java/org/apache/rocketmq/broker/coldctr/ColdDataCgCtrServiceTest.java index 7ccb3422fe3..b78b1bc33b3 100644 --- a/broker/src/test/java/org/apache/rocketmq/broker/coldctr/ColdDataCgCtrServiceTest.java +++ b/broker/src/test/java/org/apache/rocketmq/broker/coldctr/ColdDataCgCtrServiceTest.java @@ -17,6 +17,7 @@ package org.apache.rocketmq.broker.coldctr; +import java.lang.reflect.Method; import org.apache.commons.lang3.reflect.FieldUtils; import org.apache.rocketmq.broker.BrokerController; import org.apache.rocketmq.common.BrokerConfig; @@ -32,6 +33,11 @@ import java.util.concurrent.atomic.AtomicLong; 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.never; +import static org.mockito.Mockito.verify; import static org.mockito.Mockito.when; @RunWith(MockitoJUnitRunner.class) @@ -63,6 +69,25 @@ public void testGetColdDataFlowCtrInfo() { assertTrue(actual.contains("\"runtimeTable\":{\"consumerGroup1\":{\"coldAcc\":1,\"createTimeMills\":1,\"lastColdReadTimeMills\":1}}")); } + @Test + public void testDeceleratesHighestThresholdsWithoutOverflow() throws Exception { + ConcurrentHashMap thresholds = new ConcurrentHashMap<>(); + thresholds.put("highest||adaptive", Long.MAX_VALUE); + thresholds.put("medium2||adaptive", 2L); + thresholds.put("medium1||adaptive", 1L); + thresholds.put("lowest||adaptive", 0L); + ColdCtrStrategy strategy = mock(ColdCtrStrategy.class); + FieldUtils.writeField(coldDataCgCtrService, "cgColdThresholdMapConfig", thresholds, true); + FieldUtils.writeField(coldDataCgCtrService, "coldCtrStrategy", strategy, true); + Method method = ColdDataCgCtrService.class.getDeclaredMethod("sortAndDecelerate"); + method.setAccessible(true); + + method.invoke(coldDataCgCtrService); + + verify(strategy).decelerate(eq("highest||adaptive"), anyLong()); + verify(strategy, never()).decelerate(eq("lowest||adaptive"), anyLong()); + } + private Map createCgColdThresholdMapRuntime() { Map result = new ConcurrentHashMap<>(); AccAndTimeStamp accAndTimeStamp = new AccAndTimeStamp(new AtomicLong(1L));