From fe88fb36a3e1790cdb33ca670306c2adea280ba3 Mon Sep 17 00:00:00 2001 From: cserwen Date: Wed, 12 Jan 2022 21:00:56 +0800 Subject: [PATCH] [ISSUE #3498] Make messages in reviveTopic more evenly written to different queues #3499 --- .../apache/rocketmq/broker/processor/PopMessageProcessor.java | 4 +++- 1 file changed, 3 insertions(+), 1 deletion(-) diff --git a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopMessageProcessor.java b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopMessageProcessor.java index aa97fc84d0..fcc972d92c 100644 --- a/broker/src/main/java/org/apache/rocketmq/broker/processor/PopMessageProcessor.java +++ b/broker/src/main/java/org/apache/rocketmq/broker/processor/PopMessageProcessor.java @@ -92,6 +92,7 @@ public class PopMessageProcessor implements NettyRequestProcessor { private PopLongPollingService popLongPollingService; private PopBufferMergeService popBufferMergeService; private QueueLockManager queueLockManager; + private AtomicLong ckMessageNumber; public PopMessageProcessor(final BrokerController brokerController) { this.brokerController = brokerController; @@ -104,6 +105,7 @@ public class PopMessageProcessor implements NettyRequestProcessor { this.popLongPollingService = new PopLongPollingService(); this.queueLockManager = new QueueLockManager(); this.popBufferMergeService = new PopBufferMergeService(this.brokerController, this); + this.ckMessageNumber = new AtomicLong(); } public PopLongPollingService getPopLongPollingService() { @@ -350,7 +352,7 @@ public class PopMessageProcessor implements NettyRequestProcessor { if (requestHeader.isOrder()) { reviveQid = KeyBuilder.POP_ORDER_REVIVE_QUEUE; } else { - reviveQid = randomQ % this.brokerController.getBrokerConfig().getReviveQueueNum(); + reviveQid = (int) Math.abs(ckMessageNumber.getAndIncrement() % this.brokerController.getBrokerConfig().getReviveQueueNum()); } GetMessageResult getMessageResult = new GetMessageResult();