Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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());
Expand Down Expand Up @@ -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()) {
Expand Down Expand Up @@ -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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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();
}
}
}