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

lollipopjin 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 4fd0e3beea [ISSUE #11085] Stabilize offset reset tests and guard POP 
revive cache logging (#11086)
4fd0e3beea is described below

commit 4fd0e3beeafbb57b8c3050edf5cebcbfde0d2e25
Author: qianye <[email protected]>
AuthorDate: Wed Sep 9 15:15:39 2026 +0800

    [ISSUE #11085] Stabilize offset reset tests and guard POP revive cache 
logging (#11086)
---
 .../rocketmq/broker/pop/PopConsumerService.java    |   2 +-
 .../broker/pop/PopConsumerServiceTest.java         |  16 +++
 .../rocketmq/test/offset/OffsetResetForPopIT.java  | 146 +++++++++------------
 .../apache/rocketmq/test/offset/OffsetResetIT.java |  13 +-
 4 files changed, 83 insertions(+), 94 deletions(-)

diff --git 
a/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerService.java 
b/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerService.java
index f72e2ba26f..7d7aec82f8 100644
--- 
a/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerService.java
+++ 
b/broker/src/main/java/org/apache/rocketmq/broker/pop/PopConsumerService.java
@@ -646,7 +646,7 @@ public class PopConsumerService extends ServiceThread {
         currentTime.set(consumerRecords.isEmpty() ?
             upperTime : consumerRecords.get(consumerRecords.size() - 
1).getVisibilityTimeout());
 
-        if (brokerConfig.isEnablePopBufferMerge()) {
+        if (brokerConfig.isEnablePopBufferMerge() && popConsumerCache != null) 
{
             log.info("PopConsumerService, key size={}, cache size={}, revive 
count={}, failure count={}, " +
                     "behindInMillis={}, scanInMillis={}, costInMillis={}",
                 popConsumerCache.getCacheKeySize(), 
popConsumerCache.getCacheSize(),
diff --git 
a/broker/src/test/java/org/apache/rocketmq/broker/pop/PopConsumerServiceTest.java
 
b/broker/src/test/java/org/apache/rocketmq/broker/pop/PopConsumerServiceTest.java
index 44189744b4..8c61ae2787 100644
--- 
a/broker/src/test/java/org/apache/rocketmq/broker/pop/PopConsumerServiceTest.java
+++ 
b/broker/src/test/java/org/apache/rocketmq/broker/pop/PopConsumerServiceTest.java
@@ -132,6 +132,22 @@ public class PopConsumerServiceTest {
         return popConsumerRecord;
     }
 
+    @Test
+    public void reviveAfterEnablingBufferWithoutCacheTest() throws 
IllegalAccessException {
+        BrokerConfig brokerConfig = brokerController.getBrokerConfig();
+        brokerConfig.setEnablePopBufferMerge(false);
+        consumerService = new PopConsumerService(brokerController);
+        Assert.assertNull(FieldUtils.readField(consumerService, 
"popConsumerCache", true));
+        consumerService.getPopConsumerStore().start();
+        try {
+            Assert.assertEquals(0, consumerService.revive(new AtomicLong(), 
1));
+            brokerConfig.setEnablePopBufferMerge(true);
+            Assert.assertEquals(0, consumerService.revive(new AtomicLong(), 
1));
+        } finally {
+            consumerService.shutdown();
+        }
+    }
+
     @Test
     public void isPopShouldStopTest() throws IllegalAccessException {
         Assert.assertFalse(consumerService.isPopShouldStop(groupId, topicId, 
queueId));
diff --git 
a/test/src/test/java/org/apache/rocketmq/test/offset/OffsetResetForPopIT.java 
b/test/src/test/java/org/apache/rocketmq/test/offset/OffsetResetForPopIT.java
index b2092db96a..bba468b853 100644
--- 
a/test/src/test/java/org/apache/rocketmq/test/offset/OffsetResetForPopIT.java
+++ 
b/test/src/test/java/org/apache/rocketmq/test/offset/OffsetResetForPopIT.java
@@ -18,11 +18,9 @@
 package org.apache.rocketmq.test.offset;
 
 import com.google.common.collect.Lists;
-import java.util.Collections;
+import java.util.HashSet;
 import java.util.List;
 import java.util.Set;
-import java.util.concurrent.ConcurrentHashMap;
-import java.util.concurrent.Executors;
 import java.util.concurrent.TimeUnit;
 import java.util.concurrent.atomic.AtomicInteger;
 import org.apache.rocketmq.client.consumer.PopResult;
@@ -58,6 +56,7 @@ public class OffsetResetForPopIT extends BaseConf {
     private RMQNormalProducer producer = null;
     private RMQPopConsumer consumer = null;
     private DefaultMQAdminExt adminExt;
+    private boolean consumerStarted;
 
     @Before
     public void setUp() throws Exception {
@@ -77,7 +76,13 @@ public class OffsetResetForPopIT extends BaseConf {
 
     @After
     public void tearDown() {
-        shutdown();
+        try {
+            if (consumerStarted) {
+                consumer.shutdown();
+            }
+        } finally {
+            shutdown();
+        }
     }
 
     private void createAndWaitTopicRegister(String brokerName, String topic) 
throws Exception {
@@ -91,27 +96,26 @@ public class OffsetResetForPopIT extends BaseConf {
             () -> MQAdminTestUtils.checkTopicExist(adminExt, topic));
     }
 
-    private void resetOffsetInner(long resetOffset) {
-        try {
-            // reset offset by queue
-            adminExt.resetOffsetByQueueId(brokerController1.getBrokerAddr(),
-                consumer.getConsumerGroup(), consumer.getTopic(), 0, 
resetOffset);
-        } catch (Exception ignore) {
-        }
+    private void startConsumer() {
+        consumer.start();
+        consumerStarted = true;
     }
 
-    private void ackMessageSync(MessageExt messageExt) {
-        try {
-            consumer.ackAsync(brokerController1.getBrokerAddr(),
-                messageExt.getProperty(MessageConst.PROPERTY_POP_CK)).get();
-        } catch (Exception e) {
-            e.printStackTrace();
-        }
+    private void resetOffsetInner(long resetOffset) throws Exception {
+        adminExt.resetOffsetByQueueId(brokerController1.getBrokerAddr(),
+            consumer.getConsumerGroup(), consumer.getTopic(), 0, resetOffset);
+    }
+
+    private void ackMessageSync(MessageExt messageExt) throws Exception {
+        consumer.ackAsync(brokerController1.getBrokerAddr(),
+            messageExt.getProperty(MessageConst.PROPERTY_POP_CK)).get(10, 
TimeUnit.SECONDS);
     }
 
-    private void ackMessageSync(List<MessageExt> messageExtList) {
+    private void ackMessageSync(List<MessageExt> messageExtList) throws 
Exception {
         if (messageExtList != null) {
-            messageExtList.forEach(this::ackMessageSync);
+            for (MessageExt messageExt : messageExtList) {
+                ackMessageSync(messageExt);
+            }
         }
     }
 
@@ -121,7 +125,7 @@ public class OffsetResetForPopIT extends BaseConf {
         int resetOffset = 4;
         producer.send(messageCount);
         consumer = new RMQPopConsumer(NAMESRV_ADDR, topic, "*", group, new 
RMQNormalListener());
-        consumer.start();
+        startConsumer();
 
         MessageQueue mq = new MessageQueue(topic, BROKER1_NAME, 0);
         PopResult popResult = consumer.pop(brokerController1.getBrokerAddr(), 
mq);
@@ -139,7 +143,7 @@ public class OffsetResetForPopIT extends BaseConf {
         int resetOffset = 2;
         producer.send(messageCount);
         consumer = new RMQPopConsumer(NAMESRV_ADDR, topic, "*", group, new 
RMQNormalListener());
-        consumer.start();
+        startConsumer();
 
         MessageQueue mq = new MessageQueue(topic, BROKER1_NAME, 0);
         PopResult popResult1 = 
consumer.popOrderly(brokerController1.getBrokerAddr(), mq);
@@ -171,7 +175,7 @@ public class OffsetResetForPopIT extends BaseConf {
         producer.send(messageCount);
         consumer = new RMQPopConsumer(NAMESRV_ADDR, topic, "*", group, new 
RMQNormalListener());
         resetOffsetInner(resetOffset);
-        consumer.start();
+        startConsumer();
 
         MessageQueue mq = new MessageQueue(topic, BROKER1_NAME, 0);
         PopResult popResult = 
consumer.popOrderly(brokerController1.getBrokerAddr(), mq);
@@ -192,7 +196,7 @@ public class OffsetResetForPopIT extends BaseConf {
         brokerController1.getBrokerConfig().setEnablePopBufferMerge(true);
         producer.send(messageCount);
         consumer = new RMQPopConsumer(NAMESRV_ADDR, topic, "*", group, new 
RMQNormalListener());
-        consumer.start();
+        startConsumer();
 
         MessageQueue mq = new MessageQueue(topic, BROKER1_NAME, 0);
         PopResult popResult = consumer.pop(brokerController1.getBrokerAddr(), 
mq);
@@ -242,38 +246,26 @@ public class OffsetResetForPopIT extends BaseConf {
 
         MessageQueue mq = new MessageQueue(topic, BROKER1_NAME, 0);
         AtomicInteger counter = new AtomicInteger(0);
-        consumer.start();
-        Executors.newSingleThreadScheduledExecutor().execute(() -> {
-            long start = System.currentTimeMillis();
-            while (System.currentTimeMillis() - start <= 30 * 1000L) {
-                try {
-                    PopResult popResult = 
consumer.pop(brokerController1.getBrokerAddr(), mq);
-                    if (popResult == null || popResult.getMsgFoundList() == 
null) {
-                        continue;
-                    }
-
-                    int count = 
counter.addAndGet(popResult.getMsgFoundList().size());
-                    if (needAck) {
-                        ackMessageSync(popResult.getMsgFoundList());
-                    }
-                    if (count == targetCount) {
-                        for (int offset : resetOffset) {
-                            resetOffsetInner(offset);
-                        }
-                    }
-                } catch (Exception e) {
-                    e.printStackTrace();
+        startConsumer();
+        // POP, ACK and reset are sequential; keep them within the test's 
lifetime.
+        await().pollInSameThread().pollInterval(10, 
TimeUnit.MILLISECONDS).atMost(10, TimeUnit.SECONDS).until(() -> {
+            PopResult popResult = 
consumer.pop(brokerController1.getBrokerAddr(), mq);
+            if (popResult == null || popResult.getMsgFoundList() == null) {
+                return false;
+            }
+            if (needAck) {
+                ackMessageSync(popResult.getMsgFoundList());
+            }
+            int count = counter.addAndGet(popResult.getMsgFoundList().size());
+            if (count == targetCount) {
+                for (int offset : resetOffset) {
+                    resetOffsetInner(offset);
                 }
             }
-        });
-
-        await().atMost(10, TimeUnit.SECONDS).until(() -> {
-            boolean result = true;
             if (resetFuture) {
-                result = counter.get() < 10;
+                Assert.assertTrue("Reset should skip messages", count < 10);
             }
-            result &= counter.get() >= targetCount + 10 - 
resetOffset[resetOffset.length - 1];
-            return result;
+            return count >= targetCount + 10 - resetOffset[resetOffset.length 
- 1];
         });
     }
 
@@ -312,46 +304,28 @@ public class OffsetResetForPopIT extends BaseConf {
         }
         consumer = new RMQPopConsumer(NAMESRV_ADDR, topic, "*", group, new 
RMQNormalListener(), 1);
         MessageQueue mq = new MessageQueue(topic, BROKER1_NAME, 0);
-        Set<Integer> msgReceive = Collections.newSetFromMap(new 
ConcurrentHashMap<>());
+        Set<Integer> msgReceive = new HashSet<>();
         AtomicInteger counter = new AtomicInteger(0);
-        consumer.start();
+        startConsumer();
 
-        Executors.newSingleThreadScheduledExecutor().execute(() -> {
-            long start = System.currentTimeMillis();
-            while (System.currentTimeMillis() - start <= 30 * 1000L) {
-                try {
-                    PopResult popResult = 
consumer.popOrderly(brokerController1.getBrokerAddr(), mq);
-                    if (popResult == null || popResult.getMsgFoundList() == 
null) {
-                        continue;
-                    }
-                    int count = 
counter.addAndGet(popResult.getMsgFoundList().size());
-                    for (MessageExt messageExt : popResult.getMsgFoundList()) {
-                        msgReceive.add(Integer.valueOf(new 
String(messageExt.getBody())));
-                        ackMessageSync(messageExt);
-                    }
-                    if (count == targetCount) {
-                        for (int offset : resetOffset) {
-                            resetOffsetInner(offset);
-                        }
-                    }
-                } catch (Exception e) {
-                    // do nothing;
-                }
-            }
-        });
-
-        await().atMost(10, TimeUnit.SECONDS).until(() -> {
-            boolean result = true;
-            if (expectMsgReceive.size() != msgReceive.size()) {
+        await().pollInSameThread().pollInterval(10, 
TimeUnit.MILLISECONDS).atMost(10, TimeUnit.SECONDS).until(() -> {
+            PopResult popResult = 
consumer.popOrderly(brokerController1.getBrokerAddr(), mq);
+            if (popResult == null || popResult.getMsgFoundList() == null) {
                 return false;
             }
-            if (counter.get() != expectCount) {
-                return false;
+            for (MessageExt messageExt : popResult.getMsgFoundList()) {
+                msgReceive.add(Integer.valueOf(new 
String(messageExt.getBody())));
+                ackMessageSync(messageExt);
             }
-            for (Integer expectMsg : expectMsgReceive) {
-                result &= msgReceive.contains(expectMsg);
+            int count = counter.addAndGet(popResult.getMsgFoundList().size());
+            if (count == targetCount) {
+                for (int offset : resetOffset) {
+                    resetOffsetInner(offset);
+                }
             }
-            return result;
+            return count >= expectCount;
         });
+        Assert.assertEquals(expectCount, counter.get());
+        Assert.assertEquals(new HashSet<>(expectMsgReceive), msgReceive);
     }
 }
diff --git 
a/test/src/test/java/org/apache/rocketmq/test/offset/OffsetResetIT.java 
b/test/src/test/java/org/apache/rocketmq/test/offset/OffsetResetIT.java
index 150e631df8..b001c3ab7e 100644
--- a/test/src/test/java/org/apache/rocketmq/test/offset/OffsetResetIT.java
+++ b/test/src/test/java/org/apache/rocketmq/test/offset/OffsetResetIT.java
@@ -107,12 +107,9 @@ public class OffsetResetIT extends BaseConf {
 
                 Assert.assertEquals(messageQueue.getBrokerName(), 
controller.getBrokerConfig().getBrokerName());
                 long brokerOffset = 
controller.getMessageStore().getMaxOffsetInQueue(topic, 
messageQueue.getQueueId());
-                long consumerOffset = 
controller.getConsumerOffsetManager().queryOffset(
-                    consumer.getConsumerGroup(), topic, 
messageQueue.getQueueId());
                 Assert.assertEquals(brokerOffset, 
offsetWrapper.getBrokerOffset());
-                Assert.assertEquals(consumerOffset, 
offsetWrapper.getConsumerOffset());
-
-                consumerLag += brokerOffset - consumerOffset;
+                // Consumer offsets can advance after the RPC snapshot was 
taken.
+                consumerLag += offsetWrapper.getBrokerOffset() - 
offsetWrapper.getConsumerOffset();
             }
         }
         return consumerLag;
@@ -130,12 +127,13 @@ public class OffsetResetIT extends BaseConf {
         
await().pollInterval(Duration.ofSeconds(1)).atMost(Duration.ofMinutes(3)).until(
             () -> 0L == this.getConsumerLag(topic, 
consumer.getConsumerGroup()));
 
+        // Replayed messages may arrive as soon as the first broker is reset.
+        int hasConsumeBefore = listener.getMsgIndex().get();
         for (BrokerController controller : brokerControllerList) {
             defaultMQAdminExt.resetOffsetByQueueId(controller.getBrokerAddr(),
                 consumer.getConsumerGroup(), consumer.getTopic(), 3, 0);
         }
 
-        int hasConsumeBefore = listener.getMsgIndex().get();
         int expectAfterReset = brokerControllerList.size() * msgSize;
         
await().pollInterval(Duration.ofSeconds(1)).atMost(Duration.ofMinutes(3)).until(()
 -> {
             long receive = listener.getMsgIndex().get();
@@ -157,13 +155,14 @@ public class OffsetResetIT extends BaseConf {
         
await().pollInterval(Duration.ofSeconds(1)).atMost(Duration.ofMinutes(3)).until(
             () -> 0L == this.getConsumerLag(topic, 
consumer.getConsumerGroup()));
 
+        // Replayed messages may arrive as soon as the first broker is reset.
+        int hasConsumeBefore = listener.getMsgIndex().get();
         for (BrokerController controller : brokerControllerList) {
             
defaultMQAdminExt.getDefaultMQAdminExtImpl().getMqClientInstance().getMQClientAPIImpl()
                 .invokeBrokerToResetOffset(controller.getBrokerAddr(),
                     consumer.getTopic(), consumer.getConsumerGroup(), start, 
true, 3 * 1000);
         }
 
-        int hasConsumeBefore = listener.getMsgIndex().get();
         int expectAfterReset = mqs.size() * msgSize;
         
await().pollInterval(Duration.ofSeconds(1)).atMost(Duration.ofMinutes(3)).until(()
 -> {
             long receive = listener.getMsgIndex().get();

Reply via email to