diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ProcessQueue.java b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ProcessQueue.java index a0387b5c1c6..94ec4198dd2 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ProcessQueue.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/consumer/ProcessQueue.java @@ -136,7 +136,7 @@ public boolean putMessage(final List msgs) { MessageExt old = msgTreeMap.put(msg.getQueueOffset(), msg); if (null == old) { validMsgCnt++; - this.queueOffsetMax = msg.getQueueOffset(); + this.queueOffsetMax = Math.max(this.queueOffsetMax, msg.getQueueOffset()); msgSize.addAndGet(null == msg.getBody() ? 0 : msg.getBody().length); } } diff --git a/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ProcessQueueTest.java b/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ProcessQueueTest.java index a12633be1bd..89cff05f8ea 100644 --- a/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ProcessQueueTest.java +++ b/client/src/test/java/org/apache/rocketmq/client/impl/consumer/ProcessQueueTest.java @@ -80,6 +80,18 @@ public void testCachedMessageSize() { assertThat(pq.getMsgSize().get()).isEqualTo(89 * 123); } + @Test + public void testRemoveMessagesReturnsOffsetAfterHighestOutOfOrderMessage() { + ProcessQueue pq = new ProcessQueue(); + List messages = createMessageList(2); + messages.get(0).setQueueOffset(10); + messages.get(1).setQueueOffset(5); + + pq.putMessage(messages); + + assertThat(pq.removeMessage(messages)).isEqualTo(11); + } + @Test public void testContainsMessage() { ProcessQueue pq = new ProcessQueue();