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();