This is an automated email from the ASF dual-hosted git repository.
jt2594838 pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/iotdb.git
The following commit(s) were added to refs/heads/master by this push:
new ae163537785 [Subscription] Protect retained WAL by committed progress
(#18399)
ae163537785 is described below
commit ae163537785586b43c0549f35410f1d79219bcca
Author: Caideyipi <[email protected]>
AuthorDate: Tue Aug 25 09:33:22 2026 +0800
[Subscription] Protect retained WAL by committed progress (#18399)
* [Subscription] Protect retained WAL by committed progress
* fix(subscription): cache WAL retention boundary
* Keep subscription queue registration API compatible
---
.../consensus/iot/IoTConsensusServerImpl.java | 14 +-
.../subscription/SubscriptionQueueRegistry.java | 55 ++++-
.../SubscriptionWalRetentionCalculator.java | 12 +-
.../SubscriptionWalRetentionCalculatorTest.java | 101 ++++++++++
.../consensus/ConsensusPrefetchingQueue.java | 222 +++++++++++++++++++--
.../consensus/ConsensusPrefetchingQueueTest.java | 46 +++++
6 files changed, 429 insertions(+), 21 deletions(-)
diff --git
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
index 338bf30c1d9..e074e7204ee 100644
---
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
+++
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/IoTConsensusServerImpl.java
@@ -110,6 +110,7 @@ import java.util.concurrent.atomic.AtomicLong;
import java.util.concurrent.locks.Condition;
import java.util.concurrent.locks.Lock;
import java.util.concurrent.locks.ReentrantLock;
+import java.util.function.LongSupplier;
import java.util.regex.Pattern;
import static
org.apache.iotdb.commons.utils.FileUtils.humanReadableByteCountSI;
@@ -1144,8 +1145,10 @@ public class IoTConsensusServerImpl {
*/
public void registerSubscriptionQueue(
final BlockingQueue<IndexedConsensusRequest> queue,
- final SubscriptionWalRetentionPolicy retentionPolicy) {
- subscriptionQueueRegistry.register(queue, retentionPolicy);
+ final SubscriptionWalRetentionPolicy retentionPolicy,
+ final LongSupplier committedRetainedMinVersionIdSupplier) {
+ subscriptionQueueRegistry.register(
+ queue, retentionPolicy, committedRetainedMinVersionIdSupplier);
// Immediately re-evaluate the safe delete index with new subscription
awareness
checkAndUpdateSafeDeletedSearchIndex();
logger.info(
@@ -1288,8 +1291,8 @@ public class IoTConsensusServerImpl {
}
/**
- * Computes and updates the safe-to-delete WAL search index based on
replication progress and
- * subscription WAL retention policy.
+ * Computes and updates the safe-to-delete WAL search index based on
replication progress,
+ * subscription WAL retention policy, and consumer-group committed progress.
*
* <p>Because multiple subscription topics share one region WAL, the
effective per-region
* retention policy is the most conservative policy across all active
subscription queues on this
@@ -1316,7 +1319,8 @@ public class IoTConsensusServerImpl {
final SubscriptionRetentionBound subscriptionRetentionBound =
subscriptionWalRetentionCalculator.calculate(
- subscriptionQueueRegistry.getRetentionPolicies());
+ subscriptionQueueRegistry.getRetentionPolicies(),
+ subscriptionQueueRegistry.getCommittedRetainedMinVersionIds());
consensusReqReader.setSafelyDeletedSearchIndex(
Math.min(replicationIndex,
subscriptionRetentionBound.getSafelyDeletedSearchIndex()));
diff --git
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java
index a0a9422e127..8d7c27f8cf2 100644
---
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java
+++
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionQueueRegistry.java
@@ -33,6 +33,7 @@ import java.util.concurrent.BlockingQueue;
import java.util.concurrent.ConcurrentHashMap;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;
+import java.util.function.LongSupplier;
public class SubscriptionQueueRegistry {
@@ -41,7 +42,7 @@ public class SubscriptionQueueRegistry {
private static final long QUEUE_FULL_LOG_INTERVAL_MS =
TimeUnit.SECONDS.toMillis(10);
private final String consensusGroupId;
- private final Map<BlockingQueue<IndexedConsensusRequest>,
SubscriptionWalRetentionPolicy> queues =
+ private final Map<BlockingQueue<IndexedConsensusRequest>,
SubscriptionQueueRegistration> queues =
new ConcurrentHashMap<>();
private final AtomicLong droppedEntries = new AtomicLong();
private final AtomicLong lastDropLogTimeMs = new AtomicLong();
@@ -50,10 +51,26 @@ public class SubscriptionQueueRegistry {
this.consensusGroupId = consensusGroupId;
}
+ /**
+ * Registers a queue without a committed-progress constraint.
+ *
+ * <p>This overload keeps callers using the original queue-registration API
source-compatible.
+ * {@link Long#MAX_VALUE} is the neutral value for the retention calculator
and therefore does not
+ * add an extra WAL-retention constraint.
+ */
public synchronized void register(
final BlockingQueue<IndexedConsensusRequest> queue,
final SubscriptionWalRetentionPolicy retentionPolicy) {
- queues.put(queue, retentionPolicy);
+ register(queue, retentionPolicy, () -> Long.MAX_VALUE);
+ }
+
+ public synchronized void register(
+ final BlockingQueue<IndexedConsensusRequest> queue,
+ final SubscriptionWalRetentionPolicy retentionPolicy,
+ final LongSupplier committedRetainedMinVersionIdSupplier) {
+ queues.put(
+ queue,
+ new SubscriptionQueueRegistration(retentionPolicy,
committedRetainedMinVersionIdSupplier));
}
// Shares the monitor with offer() so unregister() is a real stop-receiving
barrier.
@@ -70,7 +87,26 @@ public class SubscriptionQueueRegistry {
}
public synchronized Collection<SubscriptionWalRetentionPolicy>
getRetentionPolicies() {
- return new ArrayList<>(queues.values());
+ final Collection<SubscriptionWalRetentionPolicy> retentionPolicies = new
ArrayList<>();
+ for (final SubscriptionQueueRegistration registration : queues.values()) {
+ retentionPolicies.add(registration.retentionPolicy);
+ }
+ return retentionPolicies;
+ }
+
+ public Collection<Long> getCommittedRetainedMinVersionIds() {
+ final Collection<LongSupplier> suppliers = new ArrayList<>();
+ synchronized (this) {
+ for (final SubscriptionQueueRegistration registration : queues.values())
{
+ suppliers.add(registration.committedRetainedMinVersionIdSupplier);
+ }
+ }
+
+ final Collection<Long> committedRetainedMinVersionIds = new ArrayList<>();
+ for (final LongSupplier supplier : suppliers) {
+ committedRetainedMinVersionIds.add(supplier.getAsLong());
+ }
+ return committedRetainedMinVersionIds;
}
public synchronized boolean offer(final IndexedConsensusRequest
indexedConsensusRequest) {
@@ -132,4 +168,17 @@ public class SubscriptionQueueRegistry {
}
return offeredToAnyQueue;
}
+
+ private static final class SubscriptionQueueRegistration {
+
+ private final SubscriptionWalRetentionPolicy retentionPolicy;
+ private final LongSupplier committedRetainedMinVersionIdSupplier;
+
+ private SubscriptionQueueRegistration(
+ final SubscriptionWalRetentionPolicy retentionPolicy,
+ final LongSupplier committedRetainedMinVersionIdSupplier) {
+ this.retentionPolicy = retentionPolicy;
+ this.committedRetainedMinVersionIdSupplier =
committedRetainedMinVersionIdSupplier;
+ }
+ }
}
diff --git
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculator.java
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculator.java
index 3bd2be60f37..ddddf7943de 100644
---
a/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculator.java
+++
b/iotdb-core/consensus/src/main/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculator.java
@@ -84,7 +84,8 @@ public class SubscriptionWalRetentionCalculator {
}
public SubscriptionRetentionBound calculate(
- final Collection<SubscriptionWalRetentionPolicy> retentionPolicies) {
+ final Collection<SubscriptionWalRetentionPolicy> retentionPolicies,
+ final Collection<Long> committedRetainedMinVersionIds) {
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
@@ -95,6 +96,15 @@ public class SubscriptionWalRetentionCalculator {
.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
+ // 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)));
+ }
return mergedBound;
}
diff --git
a/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculatorTest.java
b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculatorTest.java
new file mode 100644
index 00000000000..6815e8faf23
--- /dev/null
+++
b/iotdb-core/consensus/src/test/java/org/apache/iotdb/consensus/iot/subscription/SubscriptionWalRetentionCalculatorTest.java
@@ -0,0 +1,101 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements. See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership. The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License. You may obtain a copy of the License at
+ *
+ * http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied. See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.iotdb.consensus.iot.subscription;
+
+import org.apache.iotdb.consensus.iot.SubscriptionWalRetentionPolicy;
+import org.apache.iotdb.consensus.iot.log.ConsensusReqReader;
+import
org.apache.iotdb.consensus.iot.subscription.SubscriptionWalRetentionCalculator.SubscriptionRetentionBound;
+
+import org.apache.tsfile.utils.Pair;
+import org.junit.Test;
+
+import java.util.Arrays;
+import java.util.Collections;
+
+import static org.junit.Assert.assertEquals;
+
+public class SubscriptionWalRetentionCalculatorTest {
+
+ @Test
+ public void testCommittedProgressKeepsMoreWalThanTopicRetention() {
+ final SubscriptionWalRetentionCalculator calculator =
+ new SubscriptionWalRetentionCalculator(new
RetentionTestConsensusReqReader());
+ final SubscriptionRetentionBound bound =
+ calculator.calculate(
+ Collections.singletonList(
+ new SubscriptionWalRetentionPolicy(
+ "topic", 100L, SubscriptionWalRetentionPolicy.UNBOUNDED)),
+ Arrays.asList(8L, 4L));
+
+ assertEquals(100L, bound.getSafelyDeletedSearchIndex());
+ assertEquals(4L, bound.getRetainedMinVersionId());
+ }
+
+ @Test
+ public void testTopicRetentionKeepsMoreWalThanCommittedProgress() {
+ final SubscriptionWalRetentionCalculator calculator =
+ new SubscriptionWalRetentionCalculator(new
RetentionTestConsensusReqReader());
+ final SubscriptionRetentionBound bound =
+ calculator.calculate(
+ Collections.singletonList(
+ new SubscriptionWalRetentionPolicy(
+ "topic", 100L, SubscriptionWalRetentionPolicy.UNBOUNDED)),
+ Collections.singletonList(20L));
+
+ assertEquals(100L, bound.getSafelyDeletedSearchIndex());
+ assertEquals(10L, bound.getRetainedMinVersionId());
+ }
+
+ private static final class RetentionTestConsensusReqReader implements
ConsensusReqReader {
+
+ @Override
+ public void setSafelyDeletedSearchIndex(final long
safelyDeletedSearchIndex) {}
+
+ @Override
+ public ReqIterator getReqIterator(final long startIndex) {
+ throw new UnsupportedOperationException();
+ }
+
+ @Override
+ public long getCurrentSearchIndex() {
+ return 0L;
+ }
+
+ @Override
+ public long getCurrentWALFileVersion() {
+ return 0L;
+ }
+
+ @Override
+ public long getTotalSize() {
+ return 200L;
+ }
+
+ @Override
+ public long getRegionDiskUsage() {
+ return 200L;
+ }
+
+ @Override
+ public Pair<Long, Long> getDeletionBoundToFreeAtLeast(final long
bytesToFree) {
+ return new Pair<>(100L, 10L);
+ }
+ }
+}
diff --git
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
index 6bdbdbc642e..a856ac7b942 100644
---
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
+++
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueue.java
@@ -111,6 +111,21 @@ public class ConsensusPrefetchingQueue {
private final SubscriptionWalRetentionPolicy retentionPolicy;
+ /**
+ * Earliest WAL file version still needed by this consumer group's committed
progress. Topic
+ * retention and this boundary are merged conservatively by IoTConsensus
before deletion.
+ */
+ private volatile long committedRetainedMinVersionId = 0L;
+
+ private final Object committedRetentionLock = new Object();
+
+ private RegionProgress lastCommittedProgressForRetention;
+
+ private File lastRetainedWalFileForRetention;
+
+ private final Map<Long, WalFileCommitRequirement> walFileCommitRequirements =
+ new ConcurrentHashMap<>();
+
private final WakeableIndexedConsensusQueue pendingEntries;
private static final int PENDING_QUEUE_CAPACITY = 4096;
@@ -536,7 +551,8 @@ public class ConsensusPrefetchingQueue {
this.pendingEntries =
new WakeableIndexedConsensusQueue(
PENDING_QUEUE_CAPACITY, this::requestPrefetch,
this::canAcceptRealtimeEntry);
- serverImpl.registerSubscriptionQueue(pendingEntries, retentionPolicy);
+ serverImpl.registerSubscriptionQueue(
+ pendingEntries, retentionPolicy,
this::getCommittedRetainedMinVersionId);
LOGGER.info(
DataNodePipeMessages
@@ -790,6 +806,7 @@ public class ConsensusPrefetchingQueue {
this.observedSeekGeneration = seekGeneration.get();
discardBatch(this.lingerBatch);
resetBatchWriterProgress();
+ refreshCommittedWalRetentionBoundAndNotify();
LOGGER.info(
DataNodePipeMessages
@@ -1104,7 +1121,7 @@ public class ConsensusPrefetchingQueue {
compareWriterProgress(candidate, existing) > 0 ? candidate :
existing);
}
- private int compareWriterProgress(
+ private static int compareWriterProgress(
final WriterProgress leftProgress, final WriterProgress rightProgress) {
int cmp = Long.compare(leftProgress.getPhysicalTime(),
rightProgress.getPhysicalTime());
if (cmp != 0) {
@@ -2094,8 +2111,7 @@ public class ConsensusPrefetchingQueue {
commitManager.recordMapping(
consumerGroupId, topicName, consensusGroupId, writerId,
writerProgress);
if (tablets.isEmpty()) {
- return commitManager.commit(
- consumerGroupId, topicName, consensusGroupId, writerId,
writerProgress);
+ return commitAndRefreshWalRetention(writerId, writerProgress);
}
// nextOffset <= 0 means all tablets delivered in single batch
@@ -2304,8 +2320,7 @@ public class ConsensusPrefetchingQueue {
boolean committed = false;
try {
committed =
- commitManager.commitWithoutOutstanding(
- consumerGroupId, topicName, consensusGroupId, commitWriterId,
commitWriterProgress);
+ commitWithoutOutstandingAndRefreshWalRetention(commitWriterId,
commitWriterProgress);
} finally {
if (!committed
&& Objects.nonNull(event)
@@ -2710,12 +2725,7 @@ public class ConsensusPrefetchingQueue {
}
final boolean committed =
- commitManager.commit(
- consumerGroupId,
- topicName,
- consensusGroupId,
- commitWriterId,
- commitWriterProgress);
+ commitAndRefreshWalRetention(commitWriterId,
commitWriterProgress);
if (!committed) {
if (!silent) {
LOGGER.warn(
@@ -3230,6 +3240,7 @@ public class ConsensusPrefetchingQueue {
// entry so seek/rebind resumes from the intended frontier.
commitManager.resetState(
consumerGroupId, topicName, consensusGroupId,
request.committedRegionProgress);
+ refreshCommittedWalRetentionBoundAndNotify();
LOGGER.info(
DataNodePipeMessages
@@ -3243,6 +3254,122 @@ public class ConsensusPrefetchingQueue {
seekGeneration.get());
}
+ private boolean commitAndRefreshWalRetention(
+ final WriterId writerId, final WriterProgress writerProgress) {
+ final boolean committed =
+ commitManager.commit(
+ consumerGroupId, topicName, consensusGroupId, writerId,
writerProgress);
+ if (committed) {
+ refreshCommittedWalRetentionBoundAndNotify();
+ }
+ return committed;
+ }
+
+ private boolean commitWithoutOutstandingAndRefreshWalRetention(
+ final WriterId writerId, final WriterProgress writerProgress) {
+ final boolean committed =
+ commitManager.commitWithoutOutstanding(
+ consumerGroupId, topicName, consensusGroupId, writerId,
writerProgress);
+ if (committed) {
+ refreshCommittedWalRetentionBoundAndNotify();
+ }
+ return committed;
+ }
+
+ private long getCommittedRetainedMinVersionId() {
+ refreshCommittedWalRetentionBound();
+ return committedRetainedMinVersionId;
+ }
+
+ private void refreshCommittedWalRetentionBoundAndNotify() {
+ if (refreshCommittedWalRetentionBound()) {
+ serverImpl.checkAndUpdateSafeDeletedSearchIndex();
+ }
+ }
+
+ private boolean refreshCommittedWalRetentionBound() {
+ final RegionProgress committedRegionProgress =
+ commitManager.getCommittedRegionProgress(consumerGroupId, topicName,
consensusGroupId);
+
+ synchronized (committedRetentionLock) {
+ if (Objects.equals(lastCommittedProgressForRetention,
committedRegionProgress)
+ && Objects.nonNull(lastRetainedWalFileForRetention)
+ && lastRetainedWalFileForRetention.exists()) {
+ return false;
+ }
+
+ final CommittedWalRetentionBound newRetentionBound =
+ computeCommittedRetainedMinVersionId(committedRegionProgress);
+ final long newRetainedMinVersionId =
newRetentionBound.retainedMinVersionId;
+ final boolean changed = committedRetainedMinVersionId !=
newRetainedMinVersionId;
+ committedRetainedMinVersionId = newRetainedMinVersionId;
+ lastCommittedProgressForRetention = committedRegionProgress;
+ lastRetainedWalFileForRetention = newRetentionBound.retainedWalFile;
+ walFileCommitRequirements.keySet().removeIf(versionId -> versionId <
newRetainedMinVersionId);
+ return changed;
+ }
+ }
+
+ private CommittedWalRetentionBound computeCommittedRetainedMinVersionId(
+ final RegionProgress committedRegionProgress) {
+ if (!(consensusReqReader instanceof WALNode)) {
+ return new CommittedWalRetentionBound(0L, null);
+ }
+
+ final WALNode walNode = (WALNode) consensusReqReader;
+ final long currentWalVersion = walNode.getCurrentWALFileVersion();
+ final File[] walFiles =
WALFileUtils.listAllWALFiles(walNode.getLogDirectory());
+ if (Objects.isNull(walFiles) || walFiles.length == 0) {
+ return new CommittedWalRetentionBound(Math.max(0L, currentWalVersion),
null);
+ }
+
+ WALFileUtils.ascSortByVersionId(walFiles);
+ for (final File walFile : walFiles) {
+ final long versionId = WALFileUtils.parseVersionId(walFile.getName());
+ if (versionId >= currentWalVersion) {
+ return new CommittedWalRetentionBound(Math.max(0L, currentWalVersion),
walFile);
+ }
+ if (ProgressWALIterator.isHeaderOnlyWalFile(walFile)) {
+ continue;
+ }
+
+ WalFileCommitRequirement requirement =
walFileCommitRequirements.get(versionId);
+ if (Objects.isNull(requirement)) {
+ try (final ProgressWALReader reader = openProgressWALReader(walFile)) {
+ requirement =
+ WalFileCommitRequirement.fromMetadata(
+ consensusGroupId.toString(), reader.getMetaData());
+ walFileCommitRequirements.put(versionId, requirement);
+ } catch (final IOException e) {
+ LOGGER.warn(
+ DataNodePipeMessages
+
.PIPE_LOG_CONSENSUSPREFETCHINGQUEUE_FAILED_TO_READ_WAL_METADATA_FROM_A2ED50D1,
+ this,
+ walFile,
+ e);
+ return new CommittedWalRetentionBound(versionId, walFile);
+ }
+ }
+
+ if (!requirement.isCoveredBy(committedRegionProgress)) {
+ return new CommittedWalRetentionBound(versionId, walFile);
+ }
+ }
+ return new CommittedWalRetentionBound(Math.max(0L, currentWalVersion),
null);
+ }
+
+ private static final class CommittedWalRetentionBound {
+
+ private final long retainedMinVersionId;
+ private final File retainedWalFile;
+
+ private CommittedWalRetentionBound(
+ final long retainedMinVersionId, final File retainedWalFile) {
+ this.retainedMinVersionId = retainedMinVersionId;
+ this.retainedWalFile = retainedWalFile;
+ }
+ }
+
public RegionProgress computeTailRegionProgress() {
if (!(consensusReqReader instanceof WALNode)) {
return new RegionProgress(Collections.emptyMap());
@@ -3307,6 +3434,77 @@ public class ConsensusPrefetchingQueue {
}
}
+ static final class WalFileCommitRequirement {
+
+ private final Map<WriterId, WriterProgress> requiredWriterProgress;
+ private final boolean containsUnsupportedProgress;
+
+ private WalFileCommitRequirement(
+ final Map<WriterId, WriterProgress> requiredWriterProgress,
+ final boolean containsUnsupportedProgress) {
+ this.requiredWriterProgress = requiredWriterProgress;
+ this.containsUnsupportedProgress = containsUnsupportedProgress;
+ }
+
+ static WalFileCommitRequirement fromMetadata(
+ final String regionId, final WALMetaData metadata) {
+ if (Objects.isNull(metadata)) {
+ return new WalFileCommitRequirement(Collections.emptyMap(), true);
+ }
+
+ final List<Integer> buffersSize = metadata.getBuffersSize();
+ final List<Long> physicalTimes = metadata.getPhysicalTimes();
+ final List<Short> nodeIds = metadata.getNodeIds();
+ final List<Long> localSeqs = metadata.getLocalSeqs();
+ if (physicalTimes.size() < buffersSize.size()
+ || nodeIds.size() < buffersSize.size()
+ || localSeqs.size() < buffersSize.size()) {
+ return new WalFileCommitRequirement(Collections.emptyMap(), true);
+ }
+
+ final Map<WriterId, WriterProgress> requiredWriterProgress = new
LinkedHashMap<>();
+ for (int i = 0; i < buffersSize.size(); i++) {
+ final int writerNodeId = nodeIds.get(i);
+ final long physicalTime = physicalTimes.get(i);
+ final long localSeq = localSeqs.get(i);
+ if (writerNodeId < 0 && physicalTime == 0L && localSeq < 0L) {
+ // Non-search WAL entries do not participate in subscription
progress.
+ continue;
+ }
+ if (writerNodeId < 0 || physicalTime < 0L || localSeq < 0L) {
+ // Legacy or incomplete writer metadata cannot be compared safely
with RegionProgress.
+ return new WalFileCommitRequirement(Collections.emptyMap(), true);
+ }
+
+ final WriterId writerId = new WriterId(regionId, writerNodeId);
+ final WriterProgress candidateProgress = new
WriterProgress(physicalTime, localSeq);
+ requiredWriterProgress.merge(
+ writerId,
+ candidateProgress,
+ (currentProgress, candidate) ->
+ compareWriterProgress(candidate, currentProgress) > 0
+ ? candidate
+ : currentProgress);
+ }
+ return new WalFileCommitRequirement(requiredWriterProgress, false);
+ }
+
+ boolean isCoveredBy(final RegionProgress committedRegionProgress) {
+ if (containsUnsupportedProgress ||
Objects.isNull(committedRegionProgress)) {
+ return false;
+ }
+ for (final Map.Entry<WriterId, WriterProgress> entry :
requiredWriterProgress.entrySet()) {
+ final WriterProgress committedWriterProgress =
+ committedRegionProgress.getWriterPositions().get(entry.getKey());
+ if (Objects.isNull(committedWriterProgress)
+ || compareWriterProgress(committedWriterProgress,
entry.getValue()) < 0) {
+ return false;
+ }
+ }
+ return true;
+ }
+ }
+
private void mergeTailProgress(
final Map<WriterId, WriterProgress> tailProgressByWriter, final
WALMetaData metadata) {
if (Objects.isNull(metadata)) {
diff --git
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
index da73a003808..a2907745179 100644
---
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
+++
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/subscription/broker/consensus/ConsensusPrefetchingQueueTest.java
@@ -106,6 +106,52 @@ public class ConsensusPrefetchingQueueTest {
.getModifiers()));
}
+ @Test
+ public void testWalFileCommitRequirementUsesPerWriterMaximum() {
+ final String regionId = "DataRegion[1]";
+ final WALMetaData metadata = new WALMetaData();
+ metadata.add(1, -1L, 0L, 0L, -1, -1L);
+ metadata.add(1, 1L, 0L, 100L, 7, 1L);
+ metadata.add(1, 2L, 0L, 100L, 7, 2L);
+ metadata.add(1, -1L, 0L, 200L, 8, 1L);
+
+ final ConsensusPrefetchingQueue.WalFileCommitRequirement requirement =
+
ConsensusPrefetchingQueue.WalFileCommitRequirement.fromMetadata(regionId,
metadata);
+
+ assertFalse(
+ requirement.isCoveredBy(
+ new RegionProgress(
+ Collections.singletonMap(
+ new WriterId(regionId, 7), new WriterProgress(100L,
2L)))));
+ assertFalse(
+ requirement.isCoveredBy(
+ new RegionProgress(
+ java.util.Map.of(
+ new WriterId(regionId, 7),
+ new WriterProgress(100L, 1L),
+ new WriterId(regionId, 8),
+ new WriterProgress(200L, 1L)))));
+ assertTrue(
+ requirement.isCoveredBy(
+ new RegionProgress(
+ java.util.Map.of(
+ new WriterId(regionId, 7),
+ new WriterProgress(100L, 2L),
+ new WriterId(regionId, 8),
+ new WriterProgress(200L, 1L)))));
+ }
+
+ @Test
+ public void testWalFileCommitRequirementRejectsUnsupportedWriterMetadata() {
+ final WALMetaData metadata = new WALMetaData();
+ metadata.add(1, 1L, 0L, 0L, -1, 1L);
+
+ final ConsensusPrefetchingQueue.WalFileCommitRequirement requirement =
+
ConsensusPrefetchingQueue.WalFileCommitRequirement.fromMetadata("DataRegion[1]",
metadata);
+
+ assertFalse(requirement.isCoveredBy(new
RegionProgress(Collections.emptyMap())));
+ }
+
@Test
public void testReplayStartPreservesUncoveredFollowerEntries() throws
Exception {
final String originalSystemDir =
IoTDBDescriptor.getInstance().getConfig().getSystemDir();