From 940eb745acda44d5e8152dbc0307db9e0f45ba0c Mon Sep 17 00:00:00 2001 From: lizhimins <707364882@qq.com> Date: Fri, 12 Sep 2025 15:09:56 +0800 Subject: [PATCH] [ISSUE #9695] Not use pull offset when use pop orderly consume --- .../rocketmq/broker/pop/PopConsumerService.java | 12 +++++++++--- 1 file changed, 9 insertions(+), 3 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 1138ff4afe9..277a4999cf6 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 @@ -199,8 +199,14 @@ public PopConsumerContext handleGetMessageResult(PopConsumerContext context, Get return context; } - public long getPopOffset(String groupId, String topicId, int queueId, int initMode) { - long offset = this.brokerController.getConsumerOffsetManager().queryPullOffset(groupId, topicId, queueId); + public long getPopOffset(String groupId, String topicId, int queueId, int initMode, boolean fifo) { + + // For FIFO messages, the pull offset is not used. + // This preserves compatibility when switching from pull consumer to pop consumer. + long offset = fifo ? + this.brokerController.getConsumerOffsetManager().queryOffset(groupId, topicId, queueId) : + this.brokerController.getConsumerOffsetManager().queryPullOffset(groupId, topicId, queueId); + if (offset < 0L) { try { offset = this.brokerController.getPopMessageProcessor() @@ -309,7 +315,7 @@ protected CompletableFuture getMessageAsync(CompletableFutur result.addRestCount(this.getPendingFilterCount(groupId, topicId, queueId)); return CompletableFuture.completedFuture(result); } else { - final long consumeOffset = this.getPopOffset(groupId, topicId, queueId, result.getInitMode()); + final long consumeOffset = this.getPopOffset(groupId, topicId, queueId, result.getInitMode(), result.isFifo()); return getMessageAsync(clientHost, groupId, topicId, queueId, consumeOffset, remain, filter) .thenApply(getMessageResult -> handleGetMessageResult( result, getMessageResult, topicId, queueId, retryType, consumeOffset));