This is an automated email from the ASF dual-hosted git repository.

lizhimins pushed a commit to branch develop
in repository https://gitbox.apache.org/repos/asf/rocketmq.git


The following commit(s) were added to refs/heads/develop by this push:
     new eddb235c1b [ISSUE #10722] Lite Topic may not be cleaned up when 
message store timestamp is unavailable (#10723)
eddb235c1b is described below

commit eddb235c1bd8adc5c3a1125a294405526e7bb0e3
Author: Quan <[email protected]>
AuthorDate: Mon Aug 3 14:04:44 2026 +0800

    [ISSUE #10722] Lite Topic may not be cleaned up when message store 
timestamp is unavailable (#10723)
---
 .../broker/lite/AbstractLiteLifecycleManager.java  | 59 ++++++++++++++++------
 .../lite/AbstractLiteLifecycleManagerTest.java     | 13 ++++-
 .../rocketmq/test/offset/LagCalculationIT.java     |  1 +
 .../core/MessageStoreDispatcherImplTest.java       | 21 ++++----
 4 files changed, 68 insertions(+), 26 deletions(-)

diff --git 
a/broker/src/main/java/org/apache/rocketmq/broker/lite/AbstractLiteLifecycleManager.java
 
b/broker/src/main/java/org/apache/rocketmq/broker/lite/AbstractLiteLifecycleManager.java
index f7f522b833..b8ea1ef72e 100644
--- 
a/broker/src/main/java/org/apache/rocketmq/broker/lite/AbstractLiteLifecycleManager.java
+++ 
b/broker/src/main/java/org/apache/rocketmq/broker/lite/AbstractLiteLifecycleManager.java
@@ -34,6 +34,7 @@ import java.util.List;
 import java.util.Map;
 import java.util.Set;
 import java.util.concurrent.ConcurrentHashMap;
+import java.util.concurrent.TimeUnit;
 import java.util.function.Function;
 
 import static 
org.apache.rocketmq.broker.offset.ConsumerOffsetManager.TOPIC_GROUP_SEPARATOR;
@@ -44,7 +45,7 @@ import static 
org.apache.rocketmq.broker.offset.ConsumerOffsetManager.TOPIC_GROU
  */
 public abstract class AbstractLiteLifecycleManager extends ServiceThread {
     private static final Logger LOGGER = 
LoggerFactory.getLogger(LoggerName.ROCKETMQ_POP_LITE_LOGGER_NAME);
-    private static final int MAX_INVALID_SCAN_COUNT = 5;
+    static final int MAX_INVALID_SCAN_COUNT = 5;
 
     protected final BrokerController brokerController;
     protected final String brokerName;
@@ -52,7 +53,8 @@ public abstract class AbstractLiteLifecycleManager extends 
ServiceThread {
     protected MessageStore messageStore;
     protected Map<String, Integer> ttlMap = Collections.emptyMap();
     protected Map<String, Set<String>> subscriberGroupMap = 
Collections.emptyMap();
-    protected Map<String, Integer> invalidScanCountMap = new 
ConcurrentHashMap<>();
+    protected Map<String, Integer> offsetInvalidScanCountMap = new 
ConcurrentHashMap<>();
+    protected Map<String, Integer> storeTimeInvalidScanCountMap = new 
ConcurrentHashMap<>();
 
     public AbstractLiteLifecycleManager(BrokerController brokerController, 
LiteSharding liteSharding) {
         this.brokerController = brokerController;
@@ -167,20 +169,19 @@ public abstract class AbstractLiteLifecycleManager 
extends ServiceThread {
         if (!LiteUtil.isLiteTopicQueue(lmqName)) {
             return false;
         }
-        if (maxOffset <= 0) {
-            int invalidCount = invalidScanCountMap.getOrDefault(lmqName, 0) + 
1;
-            LOGGER.warn("unexpected condition, max offset <= 0, {}, {}, 
scanCount:{}", lmqName, maxOffset, invalidCount);
-            if (invalidCount > MAX_INVALID_SCAN_COUNT) { // check more times 
in case of  concurrent issue
-                invalidScanCountMap.remove(lmqName);
-                return true;
-            }
-            invalidScanCountMap.put(lmqName, invalidCount);
-            return false;
-        } else {
-            invalidScanCountMap.remove(lmqName);
+        int offsetInvalidCount = trackInvalidCount(lmqName, maxOffset <= 0, 
offsetInvalidScanCountMap);
+        if (offsetInvalidCount > 0) {
+            // check more times in case of concurrent issue
+            LOGGER.warn("unexpected condition, max offset <= 0, {}, {}, 
scanCount:{}", lmqName, maxOffset, offsetInvalidCount);
+            return offsetInvalidCount > MAX_INVALID_SCAN_COUNT;
+        }
+        long latestStoreTime = messageStore.getMessageStoreTimeStamp(lmqName, 
0, maxOffset - 1);
+        int storeTimeInvalidCount = trackInvalidCount(lmqName, latestStoreTime 
<= 0, storeTimeInvalidScanCountMap);
+        if (storeTimeInvalidCount > 0) {
+            // bypass TTL protection on purpose, but debounce against 
transient read failures
+            LOGGER.warn("latest store time <= 0, {}, {}, scanCount:{}", 
lmqName, latestStoreTime, storeTimeInvalidCount);
+            return storeTimeInvalidCount > MAX_INVALID_SCAN_COUNT;
         }
-        long latestStoreTime =
-            
this.brokerController.getMessageStore().getMessageStoreTimeStamp(lmqName, 0, 
maxOffset - 1);
         long inactiveTime = System.currentTimeMillis() - latestStoreTime;
         if (inactiveTime < brokerController.getBrokerConfig().getMinLiteTTl()) 
{
             return false;
@@ -196,7 +197,32 @@ public abstract class AbstractLiteLifecycleManager extends 
ServiceThread {
         if (hasConsumerLag(lmqName, maxOffset, latestStoreTime, parentTopic)) {
             return false;
         }
-        return inactiveTime > minutes * 60 * 1000;
+        return inactiveTime > TimeUnit.MINUTES.toMillis(minutes);
+    }
+
+    /**
+     * Track the invalid state of the given lmq: increase the count when 
invalid, reset when recovered.
+     * The counter is removed automatically once it exceeds {@link 
#MAX_INVALID_SCAN_COUNT}.
+     *
+     * @return the current invalid count, 0 means healthy (and the counter has 
been reset)
+     */
+    private int trackInvalidCount(String lmqName, boolean invalid, Map<String, 
Integer> invalidCountMap) {
+        if (!invalid) {
+            invalidCountMap.remove(lmqName);
+            return 0;
+        }
+        int invalidCount = invalidCountMap.getOrDefault(lmqName, 0) + 1;
+        if (invalidCount > MAX_INVALID_SCAN_COUNT) {
+            invalidCountMap.remove(lmqName);
+        } else {
+            invalidCountMap.put(lmqName, invalidCount);
+        }
+        return invalidCount;
+    }
+
+    private void removeInvalidCount(String lmqName) {
+        offsetInvalidScanCountMap.remove(lmqName);
+        storeTimeInvalidScanCountMap.remove(lmqName);
     }
 
     public void deleteLmq(String parentTopic, String lmqName) {
@@ -214,6 +240,7 @@ public abstract class AbstractLiteLifecycleManager extends 
ServiceThread {
             
brokerController.getLiteSubscriptionRegistry().cleanSubscription(lmqName, 
false);
             
brokerController.getConsumerOffsetManager().getPullOffsetTable().remove(
                 lmqName + TOPIC_GROUP_SEPARATOR + MixAll.TOOLS_CONSUMER_GROUP);
+            removeInvalidCount(lmqName);
             LOGGER.info("delete lmq finish. {}, sharding:{}", lmqName, 
sharding);
         } catch (Exception e) {
             LOGGER.error("delete lmq error. {}", lmqName, e);
diff --git 
a/broker/src/test/java/org/apache/rocketmq/broker/lite/AbstractLiteLifecycleManagerTest.java
 
b/broker/src/test/java/org/apache/rocketmq/broker/lite/AbstractLiteLifecycleManagerTest.java
index 5c1ab35cd3..ddc140013c 100644
--- 
a/broker/src/test/java/org/apache/rocketmq/broker/lite/AbstractLiteLifecycleManagerTest.java
+++ 
b/broker/src/test/java/org/apache/rocketmq/broker/lite/AbstractLiteLifecycleManagerTest.java
@@ -47,6 +47,7 @@ import org.mockito.Mock;
 import org.mockito.Mockito;
 import org.mockito.junit.MockitoJUnitRunner;
 
+import static 
org.apache.rocketmq.broker.lite.AbstractLiteLifecycleManager.MAX_INVALID_SCAN_COUNT;
 import static 
org.apache.rocketmq.broker.offset.ConsumerOffsetManager.TOPIC_GROUP_SEPARATOR;
 import static org.mockito.ArgumentMatchers.anyLong;
 import static org.mockito.ArgumentMatchers.anyString;
@@ -165,7 +166,17 @@ public class AbstractLiteLifecycleManagerTest {
         Assert.assertFalse(lifecycleManager.isLiteTopicExpired(PARENT_TOPIC, 
"whatever", 10L));
 
         // maxOffset invalid
-        Assert.assertFalse(lifecycleManager.isLiteTopicExpired(PARENT_TOPIC, 
EXIST_LMQ_NAME, 0L));
+        for (int i = 0; i < MAX_INVALID_SCAN_COUNT; i++) {
+            
Assert.assertFalse(lifecycleManager.isLiteTopicExpired(PARENT_TOPIC, 
EXIST_LMQ_NAME, 0L));
+        }
+        Assert.assertTrue(lifecycleManager.isLiteTopicExpired(PARENT_TOPIC, 
EXIST_LMQ_NAME, 0L));
+
+        // storeTime invalid
+        when(messageStore.getMessageStoreTimeStamp(anyString(), anyInt(), 
anyLong())).thenReturn(-1L);
+        for (int i = 0; i < MAX_INVALID_SCAN_COUNT; i++) {
+            
Assert.assertFalse(lifecycleManager.isLiteTopicExpired(PARENT_TOPIC, 
EXIST_LMQ_NAME, 100L));
+        }
+        Assert.assertTrue(lifecycleManager.isLiteTopicExpired(PARENT_TOPIC, 
EXIST_LMQ_NAME, 100L));
 
         // less than minLiteTTl
         long mockStoreTime = System.currentTimeMillis();
diff --git 
a/test/src/test/java/org/apache/rocketmq/test/offset/LagCalculationIT.java 
b/test/src/test/java/org/apache/rocketmq/test/offset/LagCalculationIT.java
index bfed96e8cd..ffcdf5c90e 100644
--- a/test/src/test/java/org/apache/rocketmq/test/offset/LagCalculationIT.java
+++ b/test/src/test/java/org/apache/rocketmq/test/offset/LagCalculationIT.java
@@ -200,6 +200,7 @@ public class LagCalculationIT extends BaseConf {
             });
             producer.send(msgMap);
         }
+        waitForFullyDispatched();
 
         // test lag estimation for tag consumer
         for (BrokerController controller : brokerControllerList) {
diff --git 
a/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/core/MessageStoreDispatcherImplTest.java
 
b/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/core/MessageStoreDispatcherImplTest.java
index 15f06d0548..290f600208 100644
--- 
a/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/core/MessageStoreDispatcherImplTest.java
+++ 
b/tieredstore/src/test/java/org/apache/rocketmq/tieredstore/core/MessageStoreDispatcherImplTest.java
@@ -146,15 +146,18 @@ public class MessageStoreDispatcherImplTest {
             new SelectMappedBufferResult(0L, buffer.asReadOnlyBuffer(), 
buffer.remaining(), null));
         dispatcher.doScheduleDispatch(flatFile, true).join();
 
-        
Awaitility.await().pollInterval(Duration.ofSeconds(1)).atMost(Duration.ofSeconds(30)).until(()
 -> {
-            List<IndexItem> resultList1 = indexService.queryAsync(
-                mq.getTopic(), "uk", 32, 0L, 
System.currentTimeMillis()).join();
-            List<IndexItem> resultList2 = indexService.queryAsync(
-                mq.getTopic(), "uk", 120, 0L, 
System.currentTimeMillis()).join();
-            Assert.assertEquals(32, resultList1.size());
-            Assert.assertEquals(100, resultList2.size());
-            return true;
-        });
+        // Index construction is submitted to the buffer commit executor 
asynchronously,
+        // so it may not be visible right after doScheduleDispatch returns. 
Use untilAsserted
+        // here, only it retries when the assertion inside fails.
+        
Awaitility.await().pollDelay(Duration.ZERO).pollInterval(Duration.ofMillis(100))
+            .atMost(Duration.ofSeconds(30)).untilAsserted(() -> {
+                List<IndexItem> resultList1 = indexService.queryAsync(
+                    mq.getTopic(), "uk", 32, 0L, 
System.currentTimeMillis()).join();
+                List<IndexItem> resultList2 = indexService.queryAsync(
+                    mq.getTopic(), "uk", 120, 0L, 
System.currentTimeMillis()).join();
+                Assert.assertEquals(32, resultList1.size());
+                Assert.assertEquals(100, resultList2.size());
+            });
 
         Assert.assertEquals(100L, flatFile.getConsumeQueueMinOffset());
         Assert.assertEquals(200L, flatFile.getConsumeQueueMaxOffset());

Reply via email to