This is an automated email from the ASF dual-hosted git repository.

jt2594838 pushed a commit to branch dev/1.3
in repository https://gitbox.apache.org/repos/asf/iotdb.git


The following commit(s) were added to refs/heads/dev/1.3 by this push:
     new cdc95b40391 Pipe: Handle hybrid meta progress indexes (#18331) (#18383)
cdc95b40391 is described below

commit cdc95b403913f9b2944944eda9566979702b1e07
Author: Caideyipi <[email protected]>
AuthorDate: Tue Aug 4 10:56:30 2026 +0800

    Pipe: Handle hybrid meta progress indexes (#18331) (#18383)
    
    * Pipe: Handle hybrid meta progress indexes
    
    * Pipe: Handle hybrid meta progress indexes
    
    * Fix getter
    
    * Pipe: Delegate progress index extraction to subclasses
---
 .../pipe/agent/task/PipeConfigNodeTaskAgent.java   |  12 +-
 .../db/pipe/agent/task/PipeDataNodeTaskAgent.java  |  10 +-
 .../processor/aggregate/AggregateProcessor.java    |  14 +--
 .../twostage/plugin/TwoStageCountProcessor.java    |  28 +++--
 .../PipeHistoricalDataRegionTsFileSource.java      |  15 ++-
 .../plugin/TwoStageCountProcessorTest.java         |  55 +++++++++
 .../TsFileResourceProgressIndexTest.java           |   9 ++
 .../commons/consensus/index/ProgressIndex.java     |  10 ++
 .../consensus/index/impl/HybridProgressIndex.java  |  25 ++++
 .../consensus/index/impl/IoTProgressIndex.java     |   9 ++
 .../consensus/index/impl/MetaProgressIndex.java    |   9 ++
 .../consensus/index/impl/MinimumProgressIndex.java |   9 ++
 .../consensus/index/impl/RecoverProgressIndex.java |   9 ++
 .../consensus/index/impl/SimpleProgressIndex.java  |   9 ++
 .../consensus/index/impl/StateProgressIndex.java   |   9 ++
 .../index/impl/TimeWindowStateProgressIndex.java   |  11 +-
 .../pipe/source/IoTDBNonDataRegionSource.java      |  47 +++++--
 .../commons/consensus/index/ProgressIndexTest.java |  82 +++++++++++++
 .../pipe/source/IoTDBNonDataRegionSourceTest.java  | 135 +++++++++++++++++++++
 19 files changed, 472 insertions(+), 35 deletions(-)

diff --git 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/agent/task/PipeConfigNodeTaskAgent.java
 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/agent/task/PipeConfigNodeTaskAgent.java
index 156a98b79d4..3fe0b1757cf 100644
--- 
a/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/agent/task/PipeConfigNodeTaskAgent.java
+++ 
b/iotdb-core/confignode/src/main/java/org/apache/iotdb/confignode/manager/pipe/agent/task/PipeConfigNodeTaskAgent.java
@@ -50,6 +50,7 @@ import java.util.Collection;
 import java.util.Collections;
 import java.util.List;
 import java.util.Map;
+import java.util.Objects;
 import java.util.Optional;
 import java.util.concurrent.atomic.AtomicLong;
 import java.util.function.Consumer;
@@ -177,10 +178,13 @@ public class PipeConfigNodeTaskAgent extends 
PipeTaskAgent {
         }
 
         final ProgressIndex progressIndex = 
groupId2TaskMetaMap.get(regionId).getProgressIndex();
-        if (progressIndex instanceof MetaProgressIndex) {
-          if (((MetaProgressIndex) progressIndex).getIndex() + 1
-              < listeningQueueNewFirstIndex.get()) {
-            listeningQueueNewFirstIndex.set(((MetaProgressIndex) 
progressIndex).getIndex() + 1);
+        final Optional<MetaProgressIndex> metaProgressIndex =
+            Objects.isNull(progressIndex)
+                ? Optional.empty()
+                : 
progressIndex.getProgressIndexByType(MetaProgressIndex.class);
+        if (metaProgressIndex.isPresent()) {
+          if (metaProgressIndex.get().getIndex() + 1 < 
listeningQueueNewFirstIndex.get()) {
+            listeningQueueNewFirstIndex.set(metaProgressIndex.get().getIndex() 
+ 1);
           }
         } else {
           // Do not clear "minimumProgressIndex"s related queues to avoid 
clearing
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
index 17504ab6557..54e8f1752ea 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java
@@ -240,11 +240,15 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent {
         }
 
         final ProgressIndex progressIndex = pipeTaskMeta.getProgressIndex();
-        if (progressIndex instanceof MetaProgressIndex) {
-          if (((MetaProgressIndex) progressIndex).getIndex() + 1
+        final Optional<MetaProgressIndex> metaProgressIndex =
+            Objects.isNull(progressIndex)
+                ? Optional.empty()
+                : 
progressIndex.getProgressIndexByType(MetaProgressIndex.class);
+        if (metaProgressIndex.isPresent()) {
+          if (metaProgressIndex.get().getIndex() + 1
               < schemaRegionId2ListeningQueueNewFirstIndex.getOrDefault(id, 
Long.MAX_VALUE)) {
             schemaRegionId2ListeningQueueNewFirstIndex.put(
-                id, ((MetaProgressIndex) progressIndex).getIndex() + 1);
+                id, metaProgressIndex.get().getIndex() + 1);
           }
         } else {
           // Do not clear "minimumProgressIndex"s related queues to avoid 
clearing
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/aggregate/AggregateProcessor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/aggregate/AggregateProcessor.java
index f12de14ebaa..6a3ad9a6159 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/aggregate/AggregateProcessor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/aggregate/AggregateProcessor.java
@@ -331,17 +331,15 @@ public class AggregateProcessor implements PipeProcessor {
 
     // Restore window state
     final ProgressIndex index = pipeTaskMeta.getProgressIndex();
-    if (index == MinimumProgressIndex.INSTANCE) {
+    if (Objects.isNull(index) || index == MinimumProgressIndex.INSTANCE) {
       return;
     }
-    if (!(index instanceof TimeWindowStateProgressIndex)) {
-      throw new PipeException(
-          String.format(
-              "The aggregate processor does not support progressIndexType %s", 
index.getType()));
-    }
-
     final TimeWindowStateProgressIndex timeWindowStateProgressIndex =
-        (TimeWindowStateProgressIndex) index;
+        
index.getProgressIndexByType(TimeWindowStateProgressIndex.class).orElse(null);
+    // A pipe altered from another processor may not have window state yet.
+    if (Objects.isNull(timeWindowStateProgressIndex)) {
+      return;
+    }
     for (final Map.Entry<String, Pair<Long, ByteBuffer>> entry :
         
timeWindowStateProgressIndex.getTimeSeries2TimestampWindowBufferPairMap().entrySet())
 {
       final AtomicReference<TimeSeriesRuntimeState> stateReference =
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/twostage/plugin/TwoStageCountProcessor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/twostage/plugin/TwoStageCountProcessor.java
index c4a3acc50b3..55a3b370a05 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/twostage/plugin/TwoStageCountProcessor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/processor/twostage/plugin/TwoStageCountProcessor.java
@@ -136,13 +136,7 @@ public class TwoStageCountProcessor implements 
PipeProcessor {
     outputSeries = parseOutputSeries(parameters);
 
     if (Objects.nonNull(pipeTaskMeta) && 
Objects.nonNull(pipeTaskMeta.getProgressIndex())) {
-      if (pipeTaskMeta.getProgressIndex() instanceof MinimumProgressIndex) {
-        pipeTaskMeta.updateProgressIndex(
-            new StateProgressIndex(Long.MIN_VALUE, new HashMap<>(), 
MinimumProgressIndex.INSTANCE));
-      }
-
-      final StateProgressIndex stateProgressIndex =
-          (StateProgressIndex) pipeTaskMeta.getProgressIndex();
+      final StateProgressIndex stateProgressIndex = 
initializeStateProgressIndex(pipeTaskMeta);
       localCommitProgressIndex.set(stateProgressIndex.getInnerProgressIndex());
       final Binary localCountState = 
stateProgressIndex.getState().get(LOCAL_COUNT_STATE_KEY);
       localCount.set(
@@ -173,6 +167,26 @@ public class TwoStageCountProcessor implements 
PipeProcessor {
             PipeProcessorConstant.PROCESSOR_OUTPUT_SERIES_KEY, 
LEGACY_PROCESSOR_OUTPUT_SERIES_KEY));
   }
 
+  static StateProgressIndex initializeStateProgressIndex(final PipeTaskMeta 
pipeTaskMeta) {
+    final ProgressIndex progressIndex = pipeTaskMeta.getProgressIndex();
+    if (progressIndex instanceof StateProgressIndex) {
+      return (StateProgressIndex) progressIndex;
+    }
+
+    final ProgressIndex updatedProgressIndex =
+        pipeTaskMeta.updateProgressIndex(
+            new StateProgressIndex(
+                Long.MIN_VALUE, Collections.emptyMap(), 
MinimumProgressIndex.INSTANCE));
+    return updatedProgressIndex
+        .getProgressIndexByType(StateProgressIndex.class)
+        .orElseThrow(
+            () ->
+                new PipeException(
+                    String.format(
+                        "Failed to initialize StateProgressIndex from progress 
index %s.",
+                        updatedProgressIndex)));
+  }
+
   @Override
   public void process(TabletInsertionEvent tabletInsertionEvent, 
EventCollector eventCollector)
       throws Exception {
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileSource.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileSource.java
index 00fd08dc2b7..6ebdebb6efd 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileSource.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/historical/PipeHistoricalDataRegionTsFileSource.java
@@ -398,7 +398,7 @@ public class PipeHistoricalDataRegionTsFileSource 
implements PipeHistoricalDataR
 
         originalResourceList.sort(
             (o1, o2) ->
-                startIndex instanceof TimeWindowStateProgressIndex
+                Objects.nonNull(getTimeWindowStateProgressIndex(startIndex))
                     ? Long.compare(o1.getFileStartTime(), 
o2.getFileStartTime())
                     : 
o1.getMaxProgressIndex().topologicalCompareTo(o2.getMaxProgressIndex()));
         pendingQueue = new ArrayDeque<>(originalResourceList);
@@ -466,9 +466,11 @@ public class PipeHistoricalDataRegionTsFileSource 
implements PipeHistoricalDataR
   }
 
   private boolean mayTsFileContainUnprocessedData(final TsFileResource 
resource) {
-    if (startIndex instanceof TimeWindowStateProgressIndex) {
+    final TimeWindowStateProgressIndex timeWindowStateProgressIndex =
+        getTimeWindowStateProgressIndex(startIndex);
+    if (Objects.nonNull(timeWindowStateProgressIndex)) {
       // The resource is closed thus the TsFileResource#getFileEndTime() is 
safe to use
-      return ((TimeWindowStateProgressIndex) startIndex).getMinTime() <= 
resource.getFileEndTime();
+      return timeWindowStateProgressIndex.getMinTime() <= 
resource.getFileEndTime();
     }
 
     if (startIndex instanceof StateProgressIndex) {
@@ -488,6 +490,13 @@ public class PipeHistoricalDataRegionTsFileSource 
implements PipeHistoricalDataR
     return false;
   }
 
+  private TimeWindowStateProgressIndex getTimeWindowStateProgressIndex(
+      final ProgressIndex progressIndex) {
+    return Objects.isNull(progressIndex)
+        ? null
+        : 
progressIndex.getProgressIndexByType(TimeWindowStateProgressIndex.class).orElse(null);
+  }
+
   private boolean mayTsFileResourceOverlappedWithPattern(final TsFileResource 
resource) {
     // Trimming to avoid unnecessary file device getter
     if (isDbNameCoveredByPattern) {
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/processor/twostage/plugin/TwoStageCountProcessorTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/processor/twostage/plugin/TwoStageCountProcessorTest.java
index 2957ffd4e3e..599ad01bb42 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/processor/twostage/plugin/TwoStageCountProcessorTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/processor/twostage/plugin/TwoStageCountProcessorTest.java
@@ -19,7 +19,14 @@
 
 package org.apache.iotdb.db.pipe.processor.twostage.plugin;
 
+import org.apache.iotdb.commons.consensus.index.ProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.HybridProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.MetaProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.SimpleProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.StateProgressIndex;
+import 
org.apache.iotdb.commons.consensus.index.impl.TimeWindowStateProgressIndex;
 import org.apache.iotdb.commons.path.PartialPath;
+import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
 import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameterValidator;
 import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
 
@@ -44,6 +51,54 @@ public class TwoStageCountProcessorTest {
     validateOutputSeries("processor.output-series", "root.db.d.s2");
   }
 
+  @Test
+  public void testInitializeStateProgressIndexFromHybridProgressIndex() {
+    final MetaProgressIndex metaProgressIndex = new MetaProgressIndex(10L);
+    final SimpleProgressIndex simpleProgressIndex = new SimpleProgressIndex(1, 
2L);
+    final ProgressIndex hybridProgressIndex =
+        new HybridProgressIndex(metaProgressIndex)
+            .updateToMinimumEqualOrIsAfterProgressIndex(simpleProgressIndex);
+    final PipeTaskMeta pipeTaskMeta = new PipeTaskMeta(hybridProgressIndex, 0);
+
+    final StateProgressIndex stateProgressIndex =
+        TwoStageCountProcessor.initializeStateProgressIndex(pipeTaskMeta);
+
+    Assert.assertSame(stateProgressIndex, pipeTaskMeta.getProgressIndex());
+    Assert.assertEquals(
+        metaProgressIndex,
+        
stateProgressIndex.getProgressIndexByType(MetaProgressIndex.class).orElse(null));
+    Assert.assertEquals(
+        simpleProgressIndex,
+        
stateProgressIndex.getProgressIndexByType(SimpleProgressIndex.class).orElse(null));
+  }
+
+  @Test
+  public void 
testInitializeStateProgressIndexFromTimeWindowStateProgressIndex() {
+    final TimeWindowStateProgressIndex timeWindowStateProgressIndex =
+        new TimeWindowStateProgressIndex(Collections.emptyMap());
+    final PipeTaskMeta pipeTaskMeta = new 
PipeTaskMeta(timeWindowStateProgressIndex, 0);
+
+    final StateProgressIndex stateProgressIndex =
+        TwoStageCountProcessor.initializeStateProgressIndex(pipeTaskMeta);
+    Assert.assertEquals(
+        timeWindowStateProgressIndex,
+        
stateProgressIndex.getProgressIndexByType(TimeWindowStateProgressIndex.class).orElse(null));
+
+    final SimpleProgressIndex simpleProgressIndex = new SimpleProgressIndex(1, 
2L);
+    final ProgressIndex updatedProgressIndex =
+        pipeTaskMeta.updateProgressIndex(
+            new StateProgressIndex(1L, Collections.emptyMap(), 
simpleProgressIndex));
+    Assert.assertTrue(updatedProgressIndex instanceof StateProgressIndex);
+    Assert.assertEquals(
+        timeWindowStateProgressIndex,
+        updatedProgressIndex
+            .getProgressIndexByType(TimeWindowStateProgressIndex.class)
+            .orElse(null));
+    Assert.assertEquals(
+        simpleProgressIndex,
+        
updatedProgressIndex.getProgressIndexByType(SimpleProgressIndex.class).orElse(null));
+  }
+
   private PartialPath parseOutputSeries(final String key, final String value) 
throws Exception {
     return TwoStageCountProcessor.parseOutputSeries(
         new PipeParameters(Collections.singletonMap(key, value)));
diff --git 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/TsFileResourceProgressIndexTest.java
 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/TsFileResourceProgressIndexTest.java
index 87b25883fb1..a0a61e3bbd3 100644
--- 
a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/TsFileResourceProgressIndexTest.java
+++ 
b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/TsFileResourceProgressIndexTest.java
@@ -54,6 +54,7 @@ import java.util.Collections;
 import java.util.HashMap;
 import java.util.List;
 import java.util.Map;
+import java.util.Optional;
 import java.util.Random;
 import java.util.stream.IntStream;
 
@@ -213,6 +214,14 @@ public class TsFileResourceProgressIndexTest {
       throw new UnsupportedOperationException("method not implemented.");
     }
 
+    @Override
+    public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+        final Class<T> progressIndexClass) {
+      return progressIndexClass.isInstance(this)
+          ? Optional.of(progressIndexClass.cast(this))
+          : Optional.empty();
+    }
+
     @Override
     public TotalOrderSumTuple getTotalOrderSumTuple() {
       return new TotalOrderSumTuple((long) val);
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/ProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/ProgressIndex.java
index 979eee0c8db..83fdac10b5c 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/ProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/ProgressIndex.java
@@ -34,6 +34,7 @@ import java.io.OutputStream;
 import java.nio.ByteBuffer;
 import java.util.List;
 import java.util.Map;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 import java.util.stream.Collectors;
 import java.util.stream.IntStream;
@@ -158,6 +159,15 @@ public abstract class ProgressIndex implements Accountable 
{
    */
   public abstract ProgressIndexType getType();
 
+  /**
+   * Extracts a progress index of the given type from this progress index.
+   *
+   * <p>{@link StateProgressIndex} and {@link HybridProgressIndex} are 
recursively unwrapped because
+   * they may contain progress indexes from other causal chains.
+   */
+  public abstract <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      Class<T> progressIndexClass);
+
   /**
    * Get the sum of the tuples of each total order relation of the {@link 
ProgressIndex}, which is
    * used for topological sorting of the {@link ProgressIndex}.
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/HybridProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/HybridProgressIndex.java
index 2c8895532dc..724a1500911 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/HybridProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/HybridProgressIndex.java
@@ -36,6 +36,7 @@ import java.util.Collections;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.Objects;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 import java.util.stream.Collectors;
 
@@ -225,6 +226,30 @@ public class HybridProgressIndex extends ProgressIndex {
     return ProgressIndexType.HYBRID_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    if (progressIndexClass.isInstance(this)) {
+      return Optional.of(progressIndexClass.cast(this));
+    }
+
+    final Map<Short, ProgressIndex> type2Index = getType2Index();
+    // Prefer a direct component over one nested in another composite progress 
index.
+    for (final ProgressIndex progressIndex : type2Index.values()) {
+      if (progressIndexClass.isInstance(progressIndex)) {
+        return Optional.of(progressIndexClass.cast(progressIndex));
+      }
+    }
+    for (final ProgressIndex progressIndex : type2Index.values()) {
+      final Optional<T> extractedProgressIndex =
+          progressIndex.getProgressIndexByType(progressIndexClass);
+      if (extractedProgressIndex.isPresent()) {
+        return extractedProgressIndex;
+      }
+    }
+    return Optional.empty();
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     lock.readLock().lock();
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/IoTProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/IoTProgressIndex.java
index 8f6a24845aa..2de6e4cb135 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/IoTProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/IoTProgressIndex.java
@@ -35,6 +35,7 @@ import java.util.Collections;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.Objects;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 
 public class IoTProgressIndex extends ProgressIndex {
@@ -197,6 +198,14 @@ public class IoTProgressIndex extends ProgressIndex {
     return ProgressIndexType.IOT_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    return progressIndexClass.isInstance(this)
+        ? Optional.of(progressIndexClass.cast(this))
+        : Optional.empty();
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     lock.readLock().lock();
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/MetaProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/MetaProgressIndex.java
index 75322152d45..80a31a7e4f3 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/MetaProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/MetaProgressIndex.java
@@ -32,6 +32,7 @@ import java.io.InputStream;
 import java.io.OutputStream;
 import java.nio.ByteBuffer;
 import java.util.Objects;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 
 public class MetaProgressIndex extends ProgressIndex {
@@ -153,6 +154,14 @@ public class MetaProgressIndex extends ProgressIndex {
     return ProgressIndexType.META_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    return progressIndexClass.isInstance(this)
+        ? Optional.of(progressIndexClass.cast(this))
+        : Optional.empty();
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     lock.readLock().lock();
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/MinimumProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/MinimumProgressIndex.java
index e22f82c9fbb..b34e97e3493 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/MinimumProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/MinimumProgressIndex.java
@@ -28,6 +28,7 @@ import java.io.IOException;
 import java.io.InputStream;
 import java.io.OutputStream;
 import java.nio.ByteBuffer;
+import java.util.Optional;
 
 public class MinimumProgressIndex extends ProgressIndex {
 
@@ -82,6 +83,14 @@ public class MinimumProgressIndex extends ProgressIndex {
     return ProgressIndexType.MINIMUM_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    return progressIndexClass.isInstance(this)
+        ? Optional.of(progressIndexClass.cast(this))
+        : Optional.empty();
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     return TOTAL_ORDER_SUM_TUPLE;
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/RecoverProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/RecoverProgressIndex.java
index 5756594abeb..f42ad50f63a 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/RecoverProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/RecoverProgressIndex.java
@@ -36,6 +36,7 @@ import java.util.Collections;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.Objects;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 import java.util.stream.Collectors;
 
@@ -202,6 +203,14 @@ public class RecoverProgressIndex extends ProgressIndex {
     return ProgressIndexType.RECOVER_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    return progressIndexClass.isInstance(this)
+        ? Optional.of(progressIndexClass.cast(this))
+        : Optional.empty();
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     lock.readLock().lock();
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/SimpleProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/SimpleProgressIndex.java
index 26d37237256..dbee0d725a5 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/SimpleProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/SimpleProgressIndex.java
@@ -32,6 +32,7 @@ import java.io.InputStream;
 import java.io.OutputStream;
 import java.nio.ByteBuffer;
 import java.util.Objects;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 
 public class SimpleProgressIndex extends ProgressIndex {
@@ -179,6 +180,14 @@ public class SimpleProgressIndex extends ProgressIndex {
     return ProgressIndexType.SIMPLE_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    return progressIndexClass.isInstance(this)
+        ? Optional.of(progressIndexClass.cast(this))
+        : Optional.empty();
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     return new TotalOrderSumTuple(memtableFlushOrderId, (long) rebootTimes);
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/StateProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/StateProgressIndex.java
index 00a3e4dce87..96a11989e7f 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/StateProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/StateProgressIndex.java
@@ -36,6 +36,7 @@ import java.nio.ByteBuffer;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.Objects;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 
 /**
@@ -196,6 +197,14 @@ public class StateProgressIndex extends ProgressIndex {
     return ProgressIndexType.STATE_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    return progressIndexClass.isInstance(this)
+        ? Optional.of(progressIndexClass.cast(this))
+        : getInnerProgressIndex().getProgressIndexByType(progressIndexClass);
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     return innerProgressIndex.getTotalOrderSumTuple();
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/TimeWindowStateProgressIndex.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/TimeWindowStateProgressIndex.java
index 02139d2058d..f98a3a1e512 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/TimeWindowStateProgressIndex.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/consensus/index/impl/TimeWindowStateProgressIndex.java
@@ -38,6 +38,7 @@ import java.util.Collections;
 import java.util.HashMap;
 import java.util.Map;
 import java.util.Objects;
+import java.util.Optional;
 import java.util.concurrent.locks.ReentrantReadWriteLock;
 import java.util.stream.Collectors;
 
@@ -209,7 +210,7 @@ public class TimeWindowStateProgressIndex extends 
ProgressIndex {
     lock.writeLock().lock();
     try {
       if (!(progressIndex instanceof TimeWindowStateProgressIndex)) {
-        return this;
+        return ProgressIndex.blendProgressIndex(this, progressIndex);
       }
 
       final TimeWindowStateProgressIndex thisTimeWindowStateProgressIndex = 
this;
@@ -239,6 +240,14 @@ public class TimeWindowStateProgressIndex extends 
ProgressIndex {
     return ProgressIndexType.TIME_WINDOW_STATE_PROGRESS_INDEX;
   }
 
+  @Override
+  public <T extends ProgressIndex> Optional<T> getProgressIndexByType(
+      final Class<T> progressIndexClass) {
+    return progressIndexClass.isInstance(this)
+        ? Optional.of(progressIndexClass.cast(this))
+        : Optional.empty();
+  }
+
   @Override
   public TotalOrderSumTuple getTotalOrderSumTuple() {
     throw new UnsupportedOperationException(
diff --git 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/source/IoTDBNonDataRegionSource.java
 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/source/IoTDBNonDataRegionSource.java
index e8fd5e54e36..7af52fc0edc 100644
--- 
a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/source/IoTDBNonDataRegionSource.java
+++ 
b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/source/IoTDBNonDataRegionSource.java
@@ -20,8 +20,8 @@
 package org.apache.iotdb.commons.pipe.source;
 
 import org.apache.iotdb.commons.consensus.index.ProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.HybridProgressIndex;
 import org.apache.iotdb.commons.consensus.index.impl.MetaProgressIndex;
-import org.apache.iotdb.commons.consensus.index.impl.MinimumProgressIndex;
 import 
org.apache.iotdb.commons.pipe.datastructure.pattern.IoTDBPipePatternOperations;
 import org.apache.iotdb.commons.pipe.datastructure.pattern.PipePattern;
 import 
org.apache.iotdb.commons.pipe.datastructure.queue.ConcurrentIterableLinkedQueue;
@@ -36,6 +36,8 @@ import org.apache.iotdb.pipe.api.event.Event;
 import org.apache.iotdb.pipe.api.exception.PipeException;
 
 import org.apache.tsfile.utils.Pair;
+import org.slf4j.Logger;
+import org.slf4j.LoggerFactory;
 
 import java.util.LinkedList;
 import java.util.List;
@@ -45,6 +47,8 @@ import java.util.concurrent.atomic.AtomicBoolean;
 
 public abstract class IoTDBNonDataRegionSource extends IoTDBSource {
 
+  private static final Logger LOGGER = 
LoggerFactory.getLogger(IoTDBNonDataRegionSource.class);
+
   protected IoTDBPipePatternOperations pipePattern;
 
   private List<PipeSnapshotEvent> historicalEvents = new LinkedList<>();
@@ -58,6 +62,7 @@ public abstract class IoTDBNonDataRegionSource extends 
IoTDBSource {
   // If the extractor is closed, it should not be started again. This is to 
avoid the case that
   // the extractor is closed and then be reused by processor.
   protected final AtomicBoolean hasBeenClosed = new AtomicBoolean(false);
+  private final AtomicBoolean hasWarnedUnexpectedHybridProgressIndex = new 
AtomicBoolean(false);
 
   protected abstract AbstractPipeListeningQueue getListeningQueue();
 
@@ -85,14 +90,15 @@ public abstract class IoTDBNonDataRegionSource extends 
IoTDBSource {
     }
 
     final ProgressIndex progressIndex = pipeTaskMeta.getProgressIndex();
+    warnIfUnexpectedHybridProgressIndex(progressIndex);
+    final MetaProgressIndex metaProgressIndex = 
extractMetaProgressIndex(progressIndex);
     final long nextIndex =
-        progressIndex instanceof MinimumProgressIndex
+        Objects.isNull(metaProgressIndex)
                 // If the index is invalid, the queue is seen as cleared 
before and thus
                 // needs snapshot re-transferring
-                || !getListeningQueue()
-                    .isGivenNextIndexValid(((MetaProgressIndex) 
progressIndex).getIndex() + 1)
+                || 
!getListeningQueue().isGivenNextIndexValid(metaProgressIndex.getIndex() + 1)
             ? getNextIndexAfterSnapshot()
-            : ((MetaProgressIndex) progressIndex).getIndex() + 1;
+            : metaProgressIndex.getIndex() + 1;
     iterator = getListeningQueue().newIterator(nextIndex);
     super.start();
   }
@@ -222,10 +228,33 @@ public abstract class IoTDBNonDataRegionSource extends 
IoTDBSource {
   //////////////////////////// APIs provided for metric framework 
////////////////////////////
 
   public long getUnTransferredEventCount() {
-    return !(pipeTaskMeta.getProgressIndex() instanceof MinimumProgressIndex)
-        ? getListeningQueue().getTailIndex()
-            - ((MetaProgressIndex) pipeTaskMeta.getProgressIndex()).getIndex()
-            - 1
+    if (Objects.isNull(pipeTaskMeta)) {
+      return 0L;
+    }
+    final ProgressIndex progressIndex = pipeTaskMeta.getProgressIndex();
+    warnIfUnexpectedHybridProgressIndex(progressIndex);
+    final MetaProgressIndex metaProgressIndex = 
extractMetaProgressIndex(progressIndex);
+    return Objects.nonNull(metaProgressIndex)
+        ? getListeningQueue().getTailIndex() - metaProgressIndex.getIndex() - 1
         : getListeningQueue().getSize() + historicalEventsCount;
   }
+
+  private static MetaProgressIndex extractMetaProgressIndex(final 
ProgressIndex progressIndex) {
+    return Objects.isNull(progressIndex)
+        ? null
+        : 
progressIndex.getProgressIndexByType(MetaProgressIndex.class).orElse(null);
+  }
+
+  private void warnIfUnexpectedHybridProgressIndex(final ProgressIndex 
progressIndex) {
+    if (Objects.nonNull(progressIndex)
+        && 
progressIndex.getProgressIndexByType(HybridProgressIndex.class).isPresent()
+        && hasWarnedUnexpectedHybridProgressIndex.compareAndSet(false, true)) {
+      LOGGER.warn(
+          "Pipe {}@{} encountered an unexpected HybridProgressIndex in {}. 
Progress index: {}.",
+          pipeName,
+          creationTime,
+          getClass().getSimpleName(),
+          progressIndex);
+    }
+  }
 }
diff --git 
a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/consensus/index/ProgressIndexTest.java
 
b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/consensus/index/ProgressIndexTest.java
new file mode 100644
index 00000000000..07ee3ef268e
--- /dev/null
+++ 
b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/consensus/index/ProgressIndexTest.java
@@ -0,0 +1,82 @@
+/*
+ * 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.commons.consensus.index;
+
+import org.apache.iotdb.commons.consensus.index.impl.HybridProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.MetaProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.SimpleProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.StateProgressIndex;
+import 
org.apache.iotdb.commons.consensus.index.impl.TimeWindowStateProgressIndex;
+
+import org.junit.Assert;
+import org.junit.Test;
+
+import java.util.Collections;
+
+public class ProgressIndexTest {
+
+  @Test
+  public void testGetProgressIndexByTypeFromStateWrappedHybridProgressIndex() {
+    final MetaProgressIndex metaProgressIndex = new MetaProgressIndex(10L);
+    final SimpleProgressIndex simpleProgressIndex = new SimpleProgressIndex(1, 
2L);
+    final ProgressIndex hybridProgressIndex =
+        new HybridProgressIndex(metaProgressIndex)
+            .updateToMinimumEqualOrIsAfterProgressIndex(simpleProgressIndex);
+    final StateProgressIndex stateProgressIndex =
+        new StateProgressIndex(1L, Collections.emptyMap(), 
hybridProgressIndex);
+
+    Assert.assertEquals(
+        metaProgressIndex,
+        
stateProgressIndex.getProgressIndexByType(MetaProgressIndex.class).orElse(null));
+    Assert.assertEquals(
+        simpleProgressIndex,
+        
stateProgressIndex.getProgressIndexByType(SimpleProgressIndex.class).orElse(null));
+    Assert.assertSame(
+        hybridProgressIndex,
+        
stateProgressIndex.getProgressIndexByType(HybridProgressIndex.class).orElse(null));
+    Assert.assertFalse(
+        
stateProgressIndex.getProgressIndexByType(TimeWindowStateProgressIndex.class).isPresent());
+  }
+
+  @Test
+  public void 
testTimeWindowStateProgressIndexBlendsWithOtherProgressIndexTypes() {
+    final TimeWindowStateProgressIndex timeWindowStateProgressIndex =
+        new TimeWindowStateProgressIndex(Collections.emptyMap());
+    final SimpleProgressIndex simpleProgressIndex = new SimpleProgressIndex(1, 
2L);
+
+    final ProgressIndex blendedProgressIndex =
+        
timeWindowStateProgressIndex.updateToMinimumEqualOrIsAfterProgressIndex(
+            simpleProgressIndex);
+    Assert.assertTrue(blendedProgressIndex instanceof HybridProgressIndex);
+    Assert.assertEquals(
+        timeWindowStateProgressIndex,
+        blendedProgressIndex
+            .getProgressIndexByType(TimeWindowStateProgressIndex.class)
+            .orElse(null));
+    Assert.assertEquals(
+        simpleProgressIndex,
+        
blendedProgressIndex.getProgressIndexByType(SimpleProgressIndex.class).orElse(null));
+
+    final ProgressIndex reverseBlendedProgressIndex =
+        simpleProgressIndex.updateToMinimumEqualOrIsAfterProgressIndex(
+            timeWindowStateProgressIndex);
+    Assert.assertEquals(blendedProgressIndex, reverseBlendedProgressIndex);
+  }
+}
diff --git 
a/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/source/IoTDBNonDataRegionSourceTest.java
 
b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/source/IoTDBNonDataRegionSourceTest.java
new file mode 100644
index 00000000000..55f10b08285
--- /dev/null
+++ 
b/iotdb-core/node-commons/src/test/java/org/apache/iotdb/commons/pipe/source/IoTDBNonDataRegionSourceTest.java
@@ -0,0 +1,135 @@
+/*
+ * 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.commons.pipe.source;
+
+import org.apache.iotdb.commons.consensus.index.ProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.HybridProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.MetaProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.SimpleProgressIndex;
+import org.apache.iotdb.commons.consensus.index.impl.StateProgressIndex;
+import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
+import 
org.apache.iotdb.commons.pipe.datastructure.queue.listening.AbstractPipeListeningQueue;
+import org.apache.iotdb.commons.pipe.event.PipeSnapshotEvent;
+import org.apache.iotdb.commons.pipe.event.PipeWritePlanEvent;
+
+import org.junit.Assert;
+import org.junit.Test;
+import org.mockito.Mockito;
+
+import java.util.Collections;
+import java.util.Optional;
+
+public class IoTDBNonDataRegionSourceTest {
+
+  @Test
+  public void testStartWithStateWrappedHybridProgressIndex() throws Exception {
+    final AbstractPipeListeningQueue listeningQueue =
+        Mockito.mock(AbstractPipeListeningQueue.class);
+    Mockito.when(listeningQueue.isGivenNextIndexValid(11L)).thenReturn(true);
+
+    final ProgressIndex hybridProgressIndex =
+        new HybridProgressIndex(new MetaProgressIndex(10L))
+            .updateToMinimumEqualOrIsAfterProgressIndex(new 
SimpleProgressIndex(1, 2L));
+    final StateProgressIndex stateProgressIndex =
+        new StateProgressIndex(1L, Collections.emptyMap(), 
hybridProgressIndex);
+    final TestNonDataRegionSource source =
+        new TestNonDataRegionSource(listeningQueue, new 
PipeTaskMeta(stateProgressIndex, 0));
+
+    source.start();
+
+    Mockito.verify(listeningQueue).newIterator(11L);
+  }
+
+  @Test
+  public void testGetUnTransferredEventCountWithHybridProgressIndex() {
+    final AbstractPipeListeningQueue listeningQueue =
+        Mockito.mock(AbstractPipeListeningQueue.class);
+    Mockito.when(listeningQueue.getTailIndex()).thenReturn(20L);
+
+    final ProgressIndex hybridProgressIndex =
+        new HybridProgressIndex(new MetaProgressIndex(10L))
+            .updateToMinimumEqualOrIsAfterProgressIndex(new 
SimpleProgressIndex(1, 2L));
+    final TestNonDataRegionSource source =
+        new TestNonDataRegionSource(listeningQueue, new 
PipeTaskMeta(hybridProgressIndex, 0));
+
+    Assert.assertEquals(9L, source.getUnTransferredEventCount());
+  }
+
+  @Test
+  public void 
testGetUnTransferredEventCountWithHybridProgressIndexWithoutMetaIndex() {
+    final AbstractPipeListeningQueue listeningQueue =
+        Mockito.mock(AbstractPipeListeningQueue.class);
+    Mockito.when(listeningQueue.getSize()).thenReturn(7L);
+
+    final TestNonDataRegionSource source =
+        new TestNonDataRegionSource(
+            listeningQueue,
+            new PipeTaskMeta(new HybridProgressIndex(new 
SimpleProgressIndex(1, 2L)), 0));
+
+    Assert.assertEquals(7L, source.getUnTransferredEventCount());
+  }
+
+  private static final class TestNonDataRegionSource extends 
IoTDBNonDataRegionSource {
+
+    private final AbstractPipeListeningQueue listeningQueue;
+
+    private TestNonDataRegionSource(
+        final AbstractPipeListeningQueue listeningQueue, final PipeTaskMeta 
pipeTaskMeta) {
+      this.listeningQueue = listeningQueue;
+      this.pipeTaskMeta = pipeTaskMeta;
+    }
+
+    @Override
+    protected AbstractPipeListeningQueue getListeningQueue() {
+      return listeningQueue;
+    }
+
+    @Override
+    protected boolean needTransferSnapshot() {
+      return false;
+    }
+
+    @Override
+    protected void triggerSnapshot() {
+      // Do nothing
+    }
+
+    @Override
+    protected long getMaxBlockingTimeMs() {
+      return 0L;
+    }
+
+    @Override
+    protected Optional<PipeWritePlanEvent> trimRealtimeEventByPipePattern(
+        final PipeWritePlanEvent event) {
+      return Optional.of(event);
+    }
+
+    @Override
+    protected boolean isTypeListened(final PipeWritePlanEvent event) {
+      return true;
+    }
+
+    @Override
+    protected void confineHistoricalEventTransferTypes(final PipeSnapshotEvent 
event) {
+      // Do nothing
+    }
+  }
+}

Reply via email to