From 774bc67520e37d842c232329f7b8a5efe033993c Mon Sep 17 00:00:00 2001 From: SSpirits Date: Wed, 7 Dec 2022 15:08:00 +0800 Subject: [PATCH] optimize LagCalculationIT --- .../rocketmq/test/util/MQAdminTestUtils.java | 2 +- .../test/offset/LagCalculationIT.java | 50 ++++++++++++------- 2 files changed, 33 insertions(+), 19 deletions(-) diff --git a/test/src/main/java/org/apache/rocketmq/test/util/MQAdminTestUtils.java b/test/src/main/java/org/apache/rocketmq/test/util/MQAdminTestUtils.java index 554289d01d..11b00a72c6 100644 --- a/test/src/main/java/org/apache/rocketmq/test/util/MQAdminTestUtils.java +++ b/test/src/main/java/org/apache/rocketmq/test/util/MQAdminTestUtils.java @@ -314,7 +314,7 @@ public class MQAdminTestUtils { public static ConsumeStats examineConsumeStats(String brokerAddr, String topic, String group) { ConsumeStats consumeStats = null; try { - consumeStats = mqAdminExt.examineConsumeStats(brokerAddr, group, topic, Long.MAX_VALUE); + consumeStats = mqAdminExt.examineConsumeStats(brokerAddr, group, topic, 3000); } catch (Exception ignored) { } return consumeStats; 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 810118b3eb..0be18a9d33 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 @@ -96,8 +96,9 @@ public class LagCalculationIT extends BaseConf { topic, mq.getQueueId()); OffsetWrapper offsetWrapper = offsetTable.get(mq); assertEquals(brokerOffset, offsetWrapper.getBrokerOffset()); - assertEquals(consumerOffset, offsetWrapper.getConsumerOffset()); - assertEquals(pullOffset, offsetWrapper.getPullOffset()); + if (offsetWrapper.getConsumerOffset() != consumerOffset || offsetWrapper.getPullOffset() != pullOffset) { + return new Pair<>(-1L, -1L); + } lag += brokerOffset - consumerOffset; pullLag += brokerOffset - pullOffset; } @@ -106,43 +107,56 @@ public class LagCalculationIT extends BaseConf { return new Pair<>(lag, pullLag); } + public void waitForFullyDispatched() { + await().atMost(5, TimeUnit.SECONDS).until(() -> { + for (BrokerController controller : brokerControllerList) { + if (controller.getMessageStore().dispatchBehindBytes() != 0) { + return false; + } + } + return true; + }); + } + @Test - public void testCalculateLag() throws InterruptedException { + public void testCalculateLag() { int msgSize = 10; List mqs = producer.getMessageQueue(); MessageQueueMsg mqMsgs = new MessageQueueMsg(mqs, msgSize); producer.send(mqMsgs.getMsgsWithMQ()); + waitForFullyDispatched(); consumer.getListener().waitForMessageConsume(producer.getAllMsgBody(), CONSUME_TIME); - // wait for updating offset - Thread.sleep(5 * 1000); + consumer.getConsumer().getDefaultMQPushConsumerImpl().persistConsumerOffset(); - Pair pair = getLag(mqs); - assertEquals(0, (long) pair.getObject1()); - assertEquals(0, (long) pair.getObject2()); + // wait for consume all msgs + await().atMost(5, TimeUnit.SECONDS).until(() -> { + Pair lag = getLag(mqs); + return lag.getObject1() == 0 && lag.getObject2() == 0; + }); blockListener.setBlock(true); consumer.clearMsg(); producer.clearMsg(); producer.send(mqMsgs.getMsgsWithMQ()); - // wait for updating offset - Thread.sleep(5 * 1000); + waitForFullyDispatched(); - pair = getLag(mqs); - assertEquals(producer.getAllMsgBody().size(), (long) pair.getObject1()); - assertEquals(0, (long) pair.getObject2()); + // wait for pull all msgs + await().atMost(5, TimeUnit.SECONDS).until(() -> { + Pair lag = getLag(mqs); + return lag.getObject1() == producer.getAllMsgBody().size() && lag.getObject2() == 0; + }); blockListener.setBlock(false); consumer.getListener().waitForMessageConsume(producer.getAllMsgBody(), CONSUME_TIME); consumer.shutdown(); producer.clearMsg(); producer.send(mqMsgs.getMsgsWithMQ()); - // wait for updating offset - Thread.sleep(5 * 1000); + waitForFullyDispatched(); - pair = getLag(mqs); - assertEquals(producer.getAllMsgBody().size(), (long) pair.getObject1()); - assertEquals(producer.getAllMsgBody().size(), (long) pair.getObject2()); + Pair lag = getLag(mqs); + assertEquals(producer.getAllMsgBody().size(), (long) lag.getObject1()); + assertEquals(producer.getAllMsgBody().size(), (long) lag.getObject2()); } @Test