Skip to content
Merged
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 @@ -3398,6 +3398,36 @@ private long getRequiredRetainedMinVersionId() {
return Math.min(getCommittedRetainedMinVersionId(), getReplayRetainedMinVersionId());
}

private static boolean isProgressAtLeast(
final RegionProgress progress, final RegionProgress previousProgress) {
if (Objects.isNull(progress) || Objects.isNull(previousProgress)) {
return false;
}
for (final Map.Entry<WriterId, WriterProgress> entry :
previousProgress.getWriterPositions().entrySet()) {
final WriterProgress writerProgress = progress.getWriterPositions().get(entry.getKey());
if (Objects.isNull(writerProgress)
|| compareWriterProgress(writerProgress, entry.getValue()) < 0) {
return false;
}
}
return true;
}

private boolean canReuseCommittedWalRetentionBound(final RegionProgress committedRegionProgress) {
if (!isProgressAtLeast(committedRegionProgress, lastCommittedProgressForRetention)
|| Objects.isNull(lastRetainedWalFileForRetention)
|| !lastRetainedWalFileForRetention.exists()
|| WALFileUtils.parseVersionId(lastRetainedWalFileForRetention.getName())
!= committedRetainedMinVersionId) {
return false;
}

final WalFileCommitRequirement requirement =
walFileCommitRequirements.get(committedRetainedMinVersionId);
return Objects.nonNull(requirement) && !requirement.isCoveredBy(committedRegionProgress);
}

private long getCommittedRetainedMinVersionId() {
refreshCommittedWalRetentionBound();
return committedRetainedMinVersionId;
Expand Down Expand Up @@ -3440,6 +3470,13 @@ private boolean refreshCommittedWalRetentionBound() {
return false;
}

// The previous boundary file remains the first uncommitted file while its cached
// requirement is still uncovered by monotonically advancing committed progress.
if (canReuseCommittedWalRetentionBound(committedRegionProgress)) {
lastCommittedProgressForRetention = committedRegionProgress;
return false;
}

final CommittedWalRetentionBound newRetentionBound =
computeCommittedRetainedMinVersionId(committedRegionProgress);
final long newRetainedMinVersionId = newRetentionBound.retainedMinVersionId;
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -437,6 +437,126 @@ public void testReplayCursorBoundsWalRetentionWhenCommittedProgressIsAhead() thr
}
}

@Test
public void testCommittedWalRetentionReusesUncoveredBoundaryRequirement() throws Exception {
final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir();
final File systemDir = temporaryFolder.newFolder("retention-boundary-cache-system");
final File walDirectory = temporaryFolder.newFolder("retention-boundary-cache-wal");
ConsensusPrefetchingQueue queue = null;
try {
IoTDBDescriptor.getInstance().getConfig().setSystemDir(systemDir.getAbsolutePath());
final File firstWal =
new File(
walDirectory,
WALFileUtils.getLogFileName(0L, 1L, WALFileStatus.CONTAINS_SEARCH_INDEX));
final File nextWal =
new File(
walDirectory,
WALFileUtils.getLogFileName(1L, 2L, WALFileStatus.CONTAINS_SEARCH_INDEX));
final File liveWal =
new File(
walDirectory,
WALFileUtils.getLogFileName(2L, 3L, WALFileStatus.CONTAINS_SEARCH_INDEX));
writeWalMetadata(firstWal, 1L, 2L, 1002L, 7);
writeWalMetadata(nextWal, 2L, 4L, 1004L, 7);
try (WALWriter ignored = new WALWriter(liveWal, WALFileVersion.V3)) {
// Keep a live successor so both data WAL files are eligible for boundary checks.
}
assertEquals(3, WALFileUtils.listAllWALFiles(walDirectory).length);
assertFalse(ProgressWALIterator.isHeaderOnlyWalFile(firstWal));

final DataRegionId regionId = new DataRegionId(1);
final WriterId writerId = new WriterId(regionId.toString(), 7);
final AtomicReference<RegionProgress> committedProgress =
new AtomicReference<>(
new RegionProgress(
Collections.singletonMap(writerId, new WriterProgress(1000L, 0L))));
final WALNode walNode = mock(WALNode.class);
when(walNode.getLogDirectory()).thenReturn(walDirectory);
when(walNode.getSortedWalFilesSnapshot()).thenReturn(new File[] {firstWal, nextWal, liveWal});
when(walNode.getCurrentWALFileVersion()).thenReturn(2L);
assertEquals(2L, walNode.getCurrentWALFileVersion());

final IoTConsensusServerImpl serverImpl = mock(IoTConsensusServerImpl.class);
when(serverImpl.getConsensusReqReader()).thenReturn(walNode);
when(serverImpl.getWriterSafeFrontierTracker()).thenReturn(new WriterSafeFrontierTracker());
final AtomicReference<LongSupplier> retentionSupplier = new AtomicReference<>();
doAnswer(
invocation -> {
retentionSupplier.set(invocation.getArgument(2));
return null;
})
.when(serverImpl)
.registerSubscriptionQueue(any(), any(), any());

final ConsensusSubscriptionCommitManager commitManager =
mock(ConsensusSubscriptionCommitManager.class);
when(commitManager.getCommittedRegionProgress("consumerGroup", "topic", regionId))
.thenAnswer(ignored -> committedProgress.get());

final AtomicInteger progressWalReaderOpenCount = new AtomicInteger();
queue =
new ConsensusPrefetchingQueue(
"consumerGroup",
"topic",
TopicConstant.ORDER_MODE_LEADER_ONLY_VALUE,
regionId,
serverImpl,
new SubscriptionWalRetentionPolicy(
"topic",
SubscriptionWalRetentionPolicy.UNBOUNDED,
SubscriptionWalRetentionPolicy.UNBOUNDED),
mock(ConsensusLogToTabletConverter.class),
commitManager,
new RegionProgress(Collections.emptyMap()),
1L,
1L,
true) {
@Override
ProgressWALReader openProgressWALReader(final File walFile) throws IOException {
progressWalReaderOpenCount.incrementAndGet();
return super.openProgressWALReader(walFile);
}
};

assertEquals(0L, retentionSupplier.get().getAsLong());
assertEquals(1, progressWalReaderOpenCount.get());

committedProgress.set(
new RegionProgress(Collections.singletonMap(writerId, new WriterProgress(1001L, 1L))));
final Method refreshRetentionBound =
ConsensusPrefetchingQueue.class.getDeclaredMethod("refreshCommittedWalRetentionBound");
refreshRetentionBound.setAccessible(true);
assertFalse((Boolean) refreshRetentionBound.invoke(queue));
assertEquals(1, progressWalReaderOpenCount.get());
assertEquals(0L, retentionSupplier.get().getAsLong());

committedProgress.set(
new RegionProgress(Collections.singletonMap(writerId, new WriterProgress(1002L, 2L))));
assertTrue((Boolean) refreshRetentionBound.invoke(queue));
assertEquals(2, progressWalReaderOpenCount.get());

final Field committedBound =
ConsensusPrefetchingQueue.class.getDeclaredField("committedRetainedMinVersionId");
committedBound.setAccessible(true);
assertEquals(1L, committedBound.getLong(queue));
// The committed boundary advanced, but the uninspected replay cursor still protects v0.
assertEquals(0L, retentionSupplier.get().getAsLong());

committedProgress.set(
new RegionProgress(Collections.singletonMap(writerId, new WriterProgress(1001L, 1L))));
assertTrue((Boolean) refreshRetentionBound.invoke(queue));
assertEquals(3, progressWalReaderOpenCount.get());
assertEquals(0L, committedBound.getLong(queue));
assertEquals(0L, retentionSupplier.get().getAsLong());
} finally {
if (queue != null) {
queue.close();
}
IoTDBDescriptor.getInstance().getConfig().setSystemDir(originalSystemDir);
}
}

@Test
public void testReplayStartPreservesUncoveredFollowerEntries() throws Exception {
final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir();
Expand Down Expand Up @@ -2639,7 +2759,9 @@ private static void writeWalMetadata(
final WALMetaData metadata = new WALMetaData();
metadata.add(1, searchIndex, 1L, physicalTime, writerNodeId, localSeq);
try (WALWriter writer = new WALWriter(walFile, WALFileVersion.V3)) {
writer.write(ByteBuffer.wrap(new byte[] {0}), metadata);
final ByteBuffer entry = ByteBuffer.allocate(1);
entry.put((byte) 0);
writer.write(entry, metadata);
}
}

Expand Down
Loading