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 @@ -1266,9 +1266,8 @@ public ConsensusReqReader getConsensusReqReader() {
public void registerSubscriptionQueue(
final BlockingQueue<IndexedConsensusRequest> queue,
final SubscriptionWalRetentionPolicy retentionPolicy,
final LongSupplier committedRetainedMinVersionIdSupplier) {
subscriptionQueueRegistry.register(
queue, retentionPolicy, committedRetainedMinVersionIdSupplier);
final LongSupplier retainedMinVersionIdSupplier) {
subscriptionQueueRegistry.register(queue, retentionPolicy, retainedMinVersionIdSupplier);
// Immediately re-evaluate the safe delete index with new subscription awareness
checkAndUpdateSafeDeletedSearchIndex();
logger.info(
Expand Down Expand Up @@ -1440,7 +1439,7 @@ public void checkAndUpdateSafeDeletedSearchIndex() {
final SubscriptionRetentionBound subscriptionRetentionBound =
subscriptionWalRetentionCalculator.calculate(
subscriptionQueueRegistry.getRetentionPolicies(),
subscriptionQueueRegistry.getCommittedRetainedMinVersionIds());
subscriptionQueueRegistry.getRetainedMinVersionIds());

consensusReqReader.setSafelyDeletedSearchIndex(
Math.min(replicationIndex, subscriptionRetentionBound.getSafelyDeletedSearchIndex()));
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -67,10 +67,9 @@ public synchronized void register(
public synchronized void register(
final BlockingQueue<IndexedConsensusRequest> queue,
final SubscriptionWalRetentionPolicy retentionPolicy,
final LongSupplier committedRetainedMinVersionIdSupplier) {
final LongSupplier retainedMinVersionIdSupplier) {
queues.put(
queue,
new SubscriptionQueueRegistration(retentionPolicy, committedRetainedMinVersionIdSupplier));
queue, new SubscriptionQueueRegistration(retentionPolicy, retainedMinVersionIdSupplier));
}

// Shares the monitor with offer() so unregister() is a real stop-receiving barrier.
Expand All @@ -94,19 +93,19 @@ public synchronized Collection<SubscriptionWalRetentionPolicy> getRetentionPolic
return retentionPolicies;
}

public Collection<Long> getCommittedRetainedMinVersionIds() {
public Collection<Long> getRetainedMinVersionIds() {
final Collection<LongSupplier> suppliers = new ArrayList<>();
synchronized (this) {
for (final SubscriptionQueueRegistration registration : queues.values()) {
suppliers.add(registration.committedRetainedMinVersionIdSupplier);
suppliers.add(registration.retainedMinVersionIdSupplier);
}
}

final Collection<Long> committedRetainedMinVersionIds = new ArrayList<>();
final Collection<Long> retainedMinVersionIds = new ArrayList<>();
for (final LongSupplier supplier : suppliers) {
committedRetainedMinVersionIds.add(supplier.getAsLong());
retainedMinVersionIds.add(supplier.getAsLong());
}
return committedRetainedMinVersionIds;
return retainedMinVersionIds;
}

public synchronized boolean offer(final IndexedConsensusRequest indexedConsensusRequest) {
Expand Down Expand Up @@ -172,13 +171,13 @@ public synchronized boolean offer(final IndexedConsensusRequest indexedConsensus
private static final class SubscriptionQueueRegistration {

private final SubscriptionWalRetentionPolicy retentionPolicy;
private final LongSupplier committedRetainedMinVersionIdSupplier;
private final LongSupplier retainedMinVersionIdSupplier;

private SubscriptionQueueRegistration(
final SubscriptionWalRetentionPolicy retentionPolicy,
final LongSupplier committedRetainedMinVersionIdSupplier) {
final LongSupplier retainedMinVersionIdSupplier) {
this.retentionPolicy = retentionPolicy;
this.committedRetainedMinVersionIdSupplier = committedRetainedMinVersionIdSupplier;
this.retainedMinVersionIdSupplier = retainedMinVersionIdSupplier;
}
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -85,7 +85,7 @@ public SubscriptionWalRetentionCalculator(final ConsensusReqReader consensusReqR

public SubscriptionRetentionBound calculate(
final Collection<SubscriptionWalRetentionPolicy> retentionPolicies,
final Collection<Long> committedRetainedMinVersionIds) {
final Collection<Long> queueRetainedMinVersionIds) {
SubscriptionRetentionBound mergedBound = SubscriptionRetentionBound.noConstraint();
for (final SubscriptionWalRetentionPolicy policy : retentionPolicies) {
// For each topic, data can be deleted once either its size retention or its time retention
Expand All @@ -96,14 +96,14 @@ public SubscriptionRetentionBound calculate(
.mergeDeleteEither(buildTimeRetentionBound(policy.getRetentionMs()));
mergedBound = mergedBound.mergeDeleteOnlyIfBoth(perQueueBound);
}
for (final long committedRetainedMinVersionId : committedRetainedMinVersionIds) {
// Topic retention is a historical replay window, while committed progress protects data
// that a consumer group has not acknowledged yet. A WAL file can be reclaimed only when
for (final long queueRetainedMinVersionId : queueRetainedMinVersionIds) {
// Topic retention is a historical replay window, while each queue also protects WAL needed
// by its committed progress and local replay cursor. A WAL file can be reclaimed only when
// both constraints allow it, so keep the more conservative (smaller) file-version bound.
mergedBound =
mergedBound.mergeDeleteOnlyIfBoth(
SubscriptionRetentionBound.of(
Long.MAX_VALUE, Math.max(0L, committedRetainedMinVersionId)));
Long.MAX_VALUE, Math.max(0L, queueRetainedMinVersionId)));
}
return mergedBound;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -1157,6 +1157,11 @@ public File getLogDirectory() {
return logDirectory;
}

public File[] getSortedWalFilesSnapshot() {
final File[] walFiles = getSortedWalFiles();
return walFiles == null ? null : Arrays.copyOf(walFiles, walFiles.length);
}

@TestOnly
File[] getCachedSortedWalFiles() {
return sortedWalFilesCache == null
Expand All @@ -1166,7 +1171,7 @@ File[] getCachedSortedWalFiles() {

@TestOnly
File[] getSortedWalFilesForTest() {
return getSortedWalFiles();
return getSortedWalFilesSnapshot();
}

/** Get the .wal file starts with the specified version id */
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -579,7 +579,7 @@ public ConsensusPrefetchingQueue(
this::requestPrefetchForRealtimeEntry,
this::canAcceptRealtimeEntry);
serverImpl.registerSubscriptionQueue(
pendingEntries, retentionPolicy, this::getCommittedRetainedMinVersionId);
pendingEntries, retentionPolicy, this::getRequiredRetainedMinVersionId);

LOGGER.info(
DataNodePipeMessages
Expand Down Expand Up @@ -3369,11 +3369,38 @@ private boolean commitWithoutOutstandingAndRefreshWalRetention(
return committed;
}

private long getRequiredRetainedMinVersionId() {
// Committed progress can advance on another consumer or replica while this queue is still
// replaying an older local search index. Keep both boundaries so WAL deletion can never pass
// the earliest entry this queue has not inspected yet.
return Math.min(getCommittedRetainedMinVersionId(), getReplayRetainedMinVersionId());
}

private long getCommittedRetainedMinVersionId() {
refreshCommittedWalRetentionBound();
return committedRetainedMinVersionId;
}

private long getReplayRetainedMinVersionId() {
if (!(consensusReqReader instanceof WALNode)) {
return 0L;
}
final WALNode walNode = (WALNode) consensusReqReader;
return findReplayRetainedMinVersionId(
walNode.getSortedWalFilesSnapshot(), nextExpectedSearchIndex.get());
}

static long findReplayRetainedMinVersionId(
final File[] walFiles, final long nextExpectedSearchIndex) {
if (Objects.isNull(walFiles) || walFiles.length == 0) {
return 0L;
}

final int replayFileIndex =
Math.max(0, WALFileUtils.binarySearchFileBySearchIndex(walFiles, nextExpectedSearchIndex));
return WALFileUtils.parseVersionId(walFiles[replayFileIndex].getName());
}

private void refreshCommittedWalRetentionBoundAndNotify() {
if (refreshCommittedWalRetentionBound()) {
serverImpl.checkAndUpdateSafeDeletedSearchIndex();
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -82,6 +82,7 @@
import java.util.concurrent.atomic.AtomicReference;
import java.util.concurrent.locks.ReentrantReadWriteLock;
import java.util.function.BooleanSupplier;
import java.util.function.LongSupplier;

import static org.junit.Assert.assertEquals;
import static org.junit.Assert.assertFalse;
Expand Down Expand Up @@ -352,6 +353,90 @@ public void testWalFileCommitRequirementRejectsUnsupportedWriterMetadata() {
assertFalse(requirement.isCoveredBy(new RegionProgress(Collections.emptyMap())));
}

@Test
public void testReplayCursorBoundsWalRetentionWhenCommittedProgressIsAhead() throws Exception {
final String originalSystemDir = IoTDBDescriptor.getInstance().getConfig().getSystemDir();
final File systemDir = temporaryFolder.newFolder("replay-cursor-retention-system");
final File walDirectory = temporaryFolder.newFolder("replay-cursor-retention-wal");
ConsensusPrefetchingQueue queue = null;
try {
final File firstWal =
new File(
walDirectory,
WALFileUtils.getLogFileName(0L, 0L, WALFileStatus.CONTAINS_SEARCH_INDEX));
final File unreadWal =
new File(
walDirectory,
WALFileUtils.getLogFileName(1L, 100L, WALFileStatus.CONTAINS_SEARCH_INDEX));
final File liveWal =
new File(
walDirectory,
WALFileUtils.getLogFileName(2L, 200L, WALFileStatus.CONTAINS_SEARCH_INDEX));
writeWalMetadata(firstWal, 1L, 1L, 1001L, 7);
writeWalMetadata(unreadWal, 101L, 101L, 1101L, 7);
try (WALWriter ignored = new WALWriter(liveWal, WALFileVersion.V3)) {
// Keep an empty live successor so both data files are eligible for deletion.
}

final DataRegionId regionId = new DataRegionId(1);
final WALNode walNode = mock(WALNode.class);
when(walNode.getLogDirectory()).thenReturn(walDirectory);
when(walNode.getSortedWalFilesSnapshot())
.thenReturn(new File[] {firstWal, unreadWal, liveWal});
when(walNode.getCurrentWALFileVersion()).thenReturn(2L);

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))
.thenReturn(
new RegionProgress(
Collections.singletonMap(
new WriterId(regionId.toString(), 7), new WriterProgress(1101L, 101L))));

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()),
101L,
1L,
true);

assertNotNull(retentionSupplier.get());
assertEquals(1L, retentionSupplier.get().getAsLong());
final Field committedBound =
ConsensusPrefetchingQueue.class.getDeclaredField("committedRetainedMinVersionId");
committedBound.setAccessible(true);
assertEquals(2L, committedBound.getLong(queue));
} 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 @@ -2464,6 +2549,20 @@ private static Tablet createWideTablet(final int columnCount, final int rowCount
return tablet;
}

private static void writeWalMetadata(
final File walFile,
final long searchIndex,
final long localSeq,
final long physicalTime,
final int writerNodeId)
throws IOException {
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);
}
}

private static IndexedConsensusRequest createRequest(final long searchIndex) {
return new IndexedConsensusRequest(
searchIndex,
Expand Down
Loading