From 0b7291b35a3a1e3517ce5403d676159c4b6500ed Mon Sep 17 00:00:00 2001 From: cserwen Date: Thu, 17 Mar 2022 21:49:02 -0500 Subject: [PATCH] [ISSUE #3503] bugfix: the consumeOffset will be set as 0 when getMessage returns null (#3504) --- .../rocketmq/broker/processor/PopMessageProcessor.java | 10 ++++++---- 1 file changed, 6 insertions(+), 4 deletions(-) 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 fcc972d92c..371628dce9 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 @@ -498,17 +498,19 @@ public class PopMessageProcessor implements NettyRequestProcessor { getMessageTmpResult = this.brokerController.getMessageStore().getMessage(requestHeader.getConsumerGroup() , topic, queueId, offset, requestHeader.getMaxMsgNums() - getMessageResult.getMessageMapedList().size(), messageFilter); + if (getMessageTmpResult == null) { + return this.brokerController.getMessageStore().getMaxOffsetInQueue(topic, queueId) - offset + restNum; + } // maybe store offset is not correct. - if (getMessageTmpResult == null - || GetMessageStatus.OFFSET_TOO_SMALL.equals(getMessageTmpResult.getStatus()) + if (GetMessageStatus.OFFSET_TOO_SMALL.equals(getMessageTmpResult.getStatus()) || GetMessageStatus.OFFSET_OVERFLOW_BADLY.equals(getMessageTmpResult.getStatus()) || GetMessageStatus.OFFSET_FOUND_NULL.equals(getMessageTmpResult.getStatus())) { // commit offset, because the offset is not correct // If offset in store is greater than cq offset, it will cause duplicate messages, // because offset in PopBuffer is not committed. POP_LOGGER.warn("Pop initial offset, because store is no correct, {}, {}->{}", - lockKey, offset, getMessageTmpResult != null ? getMessageTmpResult.getNextBeginOffset() : "null"); - offset = getMessageTmpResult != null ? getMessageTmpResult.getNextBeginOffset() : 0; + lockKey, offset, getMessageTmpResult.getNextBeginOffset()); + offset = getMessageTmpResult.getNextBeginOffset(); this.brokerController.getConsumerOffsetManager().commitOffset(channel.remoteAddress().toString(), requestHeader.getConsumerGroup(), topic, queueId, offset); getMessageTmpResult =