From 586d998a04d4486a69bc70085ea3f802dfe902aa Mon Sep 17 00:00:00 2001 From: Rui <1685901819@qq.com> Date: Sun, 2 Aug 2026 21:24:40 +0800 Subject: [PATCH] [ISSUE #10756] Fix duplicate dispatch ConsumeQueueExt leak Signed-off-by: Rui <1685901819@qq.com> --- .../apache/rocketmq/store/ConsumeQueue.java | 8 ++- .../rocketmq/store/ConsumeQueueTest.java | 58 +++++++++++++++++++ 2 files changed, 64 insertions(+), 2 deletions(-) diff --git a/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java b/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java index 0d698dacfe1..99e594cb61f 100644 --- a/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java +++ b/store/src/main/java/org/apache/rocketmq/store/ConsumeQueue.java @@ -727,7 +727,7 @@ public void putMessagePositionInfoWrapper(DispatchRequest request) { boolean canWrite = this.messageStore.getRunningFlags().isCQWriteable(); for (int i = 0; i < maxRetries && canWrite; i++) { long tagsCode = request.getTagsCode(); - if (isExtWriteEnable()) { + if (isExtWriteEnable() && !isDispatchAlreadyApplied(request.getCommitLogOffset(), request.getMsgSize())) { ConsumeQueueExt.CqExtUnit cqExtUnit = new ConsumeQueueExt.CqExtUnit(); cqExtUnit.setFilterBitMap(request.getBitMap()); cqExtUnit.setMsgStoreTime(request.getStoreTimestamp()); @@ -836,7 +836,7 @@ public void increaseQueueOffset(QueueOffsetOperator queueOffsetOperator, Message private boolean putMessagePositionInfo(final long offset, final int size, final long tagsCode, final long cqOffset) { - if (offset + size <= this.getMaxPhysicOffset()) { + if (isDispatchAlreadyApplied(offset, size)) { // During the recovery process after broker crashes, this logs will cause the scrolling of valid logs. if (messageStore.getStateMachine().getCurrentState().isAfter(MessageStoreStateMachine.MessageStoreState.RECOVER_COMMITLOG_OK) || messageStore.getMessageStoreConfig().isEnableLogConsumeQueueRepeatedlyBuildWhenRecover()) { @@ -898,6 +898,10 @@ private boolean putMessagePositionInfo(final long offset, final int size, final return false; } + private boolean isDispatchAlreadyApplied(final long offset, final int size) { + return offset + size <= this.getMaxPhysicOffset(); + } + private void fillPreBlank(final MappedFile mappedFile, final long untilWhere) { ByteBuffer byteBuffer = ByteBuffer.allocate(CQ_STORE_UNIT_SIZE); byteBuffer.putLong(0L); diff --git a/store/src/test/java/org/apache/rocketmq/store/ConsumeQueueTest.java b/store/src/test/java/org/apache/rocketmq/store/ConsumeQueueTest.java index e8e3797d021..8e6eb14cd6f 100644 --- a/store/src/test/java/org/apache/rocketmq/store/ConsumeQueueTest.java +++ b/store/src/test/java/org/apache/rocketmq/store/ConsumeQueueTest.java @@ -775,4 +775,62 @@ public void testCorrectMinOffsetAfterAllFilesDeleted() throws IOException { FileUtils.deleteQuietly(tmpDir); } } + + @Test + public void testDuplicateDispatchDoesNotLeaveConsumeQueueExtOrphan() throws IOException { + File tmpDir = Files.createTempDirectory("duplicate-cq-ext").toFile(); + MessageStoreConfig storeConfig = new MessageStoreConfig(); + storeConfig.setStorePathRootDir(tmpDir.getAbsolutePath()); + storeConfig.setMappedFileSizeConsumeQueue(4 * ConsumeQueue.CQ_STORE_UNIT_SIZE); + storeConfig.setMappedFileSizeConsumeQueueExt(10 * ConsumeQueueExt.CqExtUnit.MIN_EXT_UNIT_SIZE); + storeConfig.setEnableConsumeQueueExt(true); + DefaultMessageStore messageStore = Mockito.mock(DefaultMessageStore.class); + Mockito.when(messageStore.getMessageStoreConfig()).thenReturn(storeConfig); + Mockito.when(messageStore.getRunningFlags()).thenReturn(new RunningFlags()); + Mockito.when(messageStore.getStoreCheckpoint()).thenReturn(Mockito.mock(StoreCheckpoint.class)); + Mockito.when(messageStore.getStateMachine()).thenReturn(new MessageStoreStateMachine(null)); + + ConsumeQueue consumeQueue = new ConsumeQueue("duplicateExtTopic", 0, tmpDir.getAbsolutePath(), + storeConfig.getMappedFileSizeConsumeQueue(), messageStore); + ConsumeQueue reloadedConsumeQueue = null; + try { + DispatchRequest first = new DispatchRequest("duplicateExtTopic", 0, 0, 10, + 100, 1000, 0, null, null, 0, 0, null); + consumeQueue.putMessagePositionInfoWrapper(first); + long firstExtAddress = getRawTagsCode(consumeQueue, 0); + + consumeQueue.putMessagePositionInfoWrapper(first); + + DispatchRequest second = new DispatchRequest("duplicateExtTopic", 0, 100, 10, + 200, 2000, 1, null, null, 0, 0, null); + consumeQueue.putMessagePositionInfoWrapper(second); + long secondExtAddress = getRawTagsCode(consumeQueue, 1); + long expectedSecondExtAddress = firstExtAddress + ConsumeQueueExt.CqExtUnit.MIN_EXT_UNIT_SIZE; + consumeQueue.flush(0); + + reloadedConsumeQueue = new ConsumeQueue("duplicateExtTopic", 0, tmpDir.getAbsolutePath(), + storeConfig.getMappedFileSizeConsumeQueue(), messageStore); + Assert.assertTrue(reloadedConsumeQueue.load()); + reloadedConsumeQueue.recover(); + + Assert.assertEquals(200, reloadedConsumeQueue.getExt(expectedSecondExtAddress).getTagsCode()); + Assert.assertEquals(expectedSecondExtAddress, secondExtAddress); + } finally { + consumeQueue.destroy(); + if (reloadedConsumeQueue != null) { + reloadedConsumeQueue.destroy(); + } + FileUtils.deleteQuietly(tmpDir); + } + } + + private long getRawTagsCode(ConsumeQueue consumeQueue, long queueOffset) { + SelectMappedBufferResult result = consumeQueue.getIndexBuffer(queueOffset); + Assert.assertNotNull(result); + try { + return result.getByteBuffer().getLong(12); + } finally { + result.release(); + } + } }