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

justinchen 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 df0fd96e20d [To dev/1.3] Pipe: Do not transfer historical tsFiles when 
restarts in realtime-only mode (#15999)
df0fd96e20d is described below

commit df0fd96e20d9ac0613f4637041207c1e6342e9ba
Author: Caideyipi <[email protected]>
AuthorDate: Tue Jul 22 18:32:27 2025 +0800

    [To dev/1.3] Pipe: Do not transfer historical tsFiles when restarts in 
realtime-only mode (#15999)
    
    * Pipe: Do not transfer historical tsFiles when restarts in realtime-only 
mode (#15996)
    
    * cp
---
 .../dataregion/IoTDBDataRegionExtractor.java       | 47 +++++----------------
 .../PipeHistoricalDataRegionTsFileExtractor.java   | 49 ++++------------------
 2 files changed, 19 insertions(+), 77 deletions(-)

diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/IoTDBDataRegionExtractor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/IoTDBDataRegionExtractor.java
index 7fd7e37139e..4f80c4edd08 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/IoTDBDataRegionExtractor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/IoTDBDataRegionExtractor.java
@@ -20,7 +20,6 @@
 package org.apache.iotdb.db.pipe.extractor.dataregion;
 
 import org.apache.iotdb.commons.consensus.DataRegionId;
-import org.apache.iotdb.commons.pipe.config.constant.SystemConstant;
 import org.apache.iotdb.commons.pipe.datastructure.pattern.IoTDBPipePattern;
 import org.apache.iotdb.commons.pipe.datastructure.pattern.PipePattern;
 import org.apache.iotdb.commons.pipe.extractor.IoTDBExtractor;
@@ -50,8 +49,6 @@ import org.apache.tsfile.utils.Pair;
 import org.slf4j.Logger;
 import org.slf4j.LoggerFactory;
 
-import javax.annotation.Nullable;
-
 import java.util.Arrays;
 import java.util.Objects;
 import java.util.concurrent.atomic.AtomicReference;
@@ -95,7 +92,7 @@ public class IoTDBDataRegionExtractor extends IoTDBExtractor {
 
   private static final Logger LOGGER = 
LoggerFactory.getLogger(IoTDBDataRegionExtractor.class);
 
-  private @Nullable PipeHistoricalDataRegionExtractor historicalExtractor;
+  private PipeHistoricalDataRegionExtractor historicalExtractor;
   private PipeRealtimeDataRegionExtractor realtimeExtractor;
 
   private DataRegionWatermarkInjector watermarkInjector;
@@ -213,22 +210,10 @@ public class IoTDBDataRegionExtractor extends 
IoTDBExtractor {
           EXTRACTOR_HISTORY_END_TIME_KEY);
     }
 
-    if (validator
-            .getParameters()
-            .getBooleanOrDefault(SystemConstant.RESTART_KEY, 
SystemConstant.RESTART_DEFAULT_VALUE)
-        || validator
-            .getParameters()
-            .getBooleanOrDefault(
-                Arrays.asList(EXTRACTOR_HISTORY_ENABLE_KEY, 
SOURCE_HISTORY_ENABLE_KEY),
-                EXTRACTOR_HISTORY_ENABLE_DEFAULT_VALUE)) {
-      // Do not flush or open historical extractor when historical tsFile is 
disabled
-      constructHistoricalExtractor();
-    }
+    constructHistoricalExtractor();
     constructRealtimeExtractor(validator.getParameters());
 
-    if (Objects.nonNull(historicalExtractor)) {
-      historicalExtractor.validate(validator);
-    }
+    historicalExtractor.validate(validator);
     realtimeExtractor.validate(validator);
   }
 
@@ -319,9 +304,7 @@ public class IoTDBDataRegionExtractor extends 
IoTDBExtractor {
 
     super.customize(parameters, configuration);
 
-    if (Objects.nonNull(historicalExtractor)) {
-      historicalExtractor.customize(parameters, configuration);
-    }
+    historicalExtractor.customize(parameters, configuration);
     realtimeExtractor.customize(parameters, configuration);
 
     // Set watermark injector
@@ -358,9 +341,7 @@ public class IoTDBDataRegionExtractor extends 
IoTDBExtractor {
         "Pipe {}@{}: Starting historical extractor {} and realtime extractor 
{}.",
         pipeName,
         regionId,
-        Objects.nonNull(historicalExtractor)
-            ? historicalExtractor.getClass().getSimpleName()
-            : null,
+        historicalExtractor.getClass().getSimpleName(),
         realtimeExtractor.getClass().getSimpleName());
 
     super.start();
@@ -395,9 +376,7 @@ public class IoTDBDataRegionExtractor extends 
IoTDBExtractor {
             "Pipe {}@{}: Started historical extractor {} and realtime 
extractor {} successfully within {} ms.",
             pipeName,
             regionId,
-            Objects.nonNull(historicalExtractor)
-                ? historicalExtractor.getClass().getSimpleName()
-                : null,
+            historicalExtractor.getClass().getSimpleName(),
             realtimeExtractor.getClass().getSimpleName(),
             System.currentTimeMillis() - startTime);
         return;
@@ -415,18 +394,14 @@ public class IoTDBDataRegionExtractor extends 
IoTDBExtractor {
       // There can still be writing when tsFile events are added. If we start
       // realtimeExtractor after the process, then this part of data will be 
lost.
       realtimeExtractor.start();
-      if (Objects.nonNull(historicalExtractor)) {
-        historicalExtractor.start();
-      }
+      historicalExtractor.start();
     } catch (final Exception e) {
       exceptionHolder.set(e);
       LOGGER.warn(
           "Pipe {}@{}: Start historical extractor {} and realtime extractor {} 
error.",
           pipeName,
           regionId,
-          Objects.nonNull(historicalExtractor)
-              ? historicalExtractor.getClass().getSimpleName()
-              : null,
+          historicalExtractor.getClass().getSimpleName(),
           realtimeExtractor.getClass().getSimpleName(),
           e);
     }
@@ -445,7 +420,7 @@ public class IoTDBDataRegionExtractor extends 
IoTDBExtractor {
     }
 
     Event event = null;
-    if (Objects.nonNull(historicalExtractor) && 
!historicalExtractor.hasConsumedAll()) {
+    if (!historicalExtractor.hasConsumedAll()) {
       event = historicalExtractor.supply();
     } else {
       if (Objects.nonNull(watermarkInjector)) {
@@ -475,9 +450,7 @@ public class IoTDBDataRegionExtractor extends 
IoTDBExtractor {
       return;
     }
 
-    if (Objects.nonNull(historicalExtractor)) {
-      historicalExtractor.close();
-    }
+    historicalExtractor.close();
     realtimeExtractor.close();
     if (Objects.nonNull(taskID)) {
       PipeDataRegionExtractorMetrics.getInstance().deregister(taskID);
diff --git 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/historical/PipeHistoricalDataRegionTsFileExtractor.java
 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/historical/PipeHistoricalDataRegionTsFileExtractor.java
index 4b6b11efed0..090b878e586 100644
--- 
a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/historical/PipeHistoricalDataRegionTsFileExtractor.java
+++ 
b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/extractor/dataregion/historical/PipeHistoricalDataRegionTsFileExtractor.java
@@ -39,7 +39,6 @@ import 
org.apache.iotdb.db.storageengine.dataregion.DataRegion;
 import org.apache.iotdb.db.storageengine.dataregion.memtable.TsFileProcessor;
 import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileManager;
 import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
-import 
org.apache.iotdb.db.storageengine.dataregion.tsfile.generator.TsFileNameGenerator;
 import org.apache.iotdb.db.utils.DateTimeUtils;
 import 
org.apache.iotdb.pipe.api.customizer.configuration.PipeExtractorRuntimeConfiguration;
 import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameterValidator;
@@ -110,7 +109,6 @@ public class PipeHistoricalDataRegionTsFileExtractor 
implements PipeHistoricalDa
   private boolean isHistoricalExtractorEnabled = false;
   private long historicalDataExtractionStartTime = Long.MIN_VALUE; // Event 
time
   private long historicalDataExtractionEndTime = Long.MAX_VALUE; // Event time
-  private long historicalDataExtractionTimeLowerBound; // Arrival time
 
   private boolean sloppyTimeRange; // true to disable time range filter after 
extraction
   private boolean sloppyPattern; // true to disable pattern filter after 
extraction
@@ -278,18 +276,6 @@ public class PipeHistoricalDataRegionTsFileExtractor 
implements PipeHistoricalDa
       }
     }
 
-    // Enable historical extractor by default
-    historicalDataExtractionTimeLowerBound =
-        isHistoricalExtractorEnabled
-            ? Long.MIN_VALUE
-            // We define the realtime data as the data generated after the 
creation time
-            // of the pipe from user's perspective. But we still need to use
-            // PipeHistoricalDataRegionExtractor to extract the realtime data 
generated between the
-            // creation time of the pipe and the time when the pipe starts, 
because those data
-            // can not be listened by PipeRealtimeDataRegionExtractor, and 
should be extracted by
-            // PipeHistoricalDataRegionExtractor from implementation 
perspective.
-            : environment.getCreationTime();
-
     shouldTransferModFile =
         parameters.getBooleanOrDefault(
             Arrays.asList(SOURCE_MODS_ENABLE_KEY, EXTRACTOR_MODS_ENABLE_KEY),
@@ -321,7 +307,7 @@ public class PipeHistoricalDataRegionTsFileExtractor 
implements PipeHistoricalDa
 
   @Override
   public synchronized void start() {
-    if (!shouldExtractInsertion || !isHistoricalExtractorEnabled) {
+    if (!shouldExtractInsertion) {
       hasBeenStarted = true;
       return;
     }
@@ -383,8 +369,10 @@ public class PipeHistoricalDataRegionTsFileExtractor 
implements PipeHistoricalDa
                 .peek(originalResourceList::add)
                 .filter(
                     resource ->
-                        // Some resource is marked as deleted but not removed 
from the list.
-                        !resource.isDeleted()
+                        isHistoricalExtractorEnabled
+                            &&
+                            // Some resource is marked as deleted but not 
removed from the list.
+                            !resource.isDeleted()
                             // Some resource is generated by pipe. We ignore 
them if the pipe should
                             // not transfer pipe requests.
                             && (!resource.isGeneratedByPipe() || 
isForwardingPipeRequests)
@@ -398,7 +386,6 @@ public class PipeHistoricalDataRegionTsFileExtractor 
implements PipeHistoricalDa
                                         .orElse(true)
                                 || mayTsFileContainUnprocessedData(resource)
                                     && 
isTsFileResourceOverlappedWithTimeRange(resource)
-                                    && 
isTsFileGeneratedAfterExtractionTimeLowerBound(resource)
                                     && 
mayTsFileResourceOverlappedWithPattern(resource)))
                 .collect(Collectors.toList());
         filteredTsFileResources.addAll(sequenceTsFileResources);
@@ -408,8 +395,10 @@ public class PipeHistoricalDataRegionTsFileExtractor 
implements PipeHistoricalDa
                 .peek(originalResourceList::add)
                 .filter(
                     resource ->
-                        // Some resource is marked as deleted but not removed 
from the list.
-                        !resource.isDeleted()
+                        isHistoricalExtractorEnabled
+                            &&
+                            // Some resource is marked as deleted but not 
removed from the list.
+                            !resource.isDeleted()
                             // Some resource is generated by pipe. We ignore 
them if the pipe should
                             // not transfer pipe requests.
                             && (!resource.isGeneratedByPipe() || 
isForwardingPipeRequests)
@@ -423,7 +412,6 @@ public class PipeHistoricalDataRegionTsFileExtractor 
implements PipeHistoricalDa
                                         .orElse(true)
                                 || mayTsFileContainUnprocessedData(resource)
                                     && 
isTsFileResourceOverlappedWithTimeRange(resource)
-                                    && 
isTsFileGeneratedAfterExtractionTimeLowerBound(resource)
                                     && 
mayTsFileResourceOverlappedWithPattern(resource)))
                 .collect(Collectors.toList());
         filteredTsFileResources.addAll(unSequenceTsFileResources);
@@ -526,25 +514,6 @@ public class PipeHistoricalDataRegionTsFileExtractor 
implements PipeHistoricalDa
         && historicalDataExtractionEndTime >= resource.getFileEndTime();
   }
 
-  private boolean isTsFileGeneratedAfterExtractionTimeLowerBound(final 
TsFileResource resource) {
-    try {
-      return historicalDataExtractionTimeLowerBound
-          <= 
TsFileNameGenerator.getTsFileName(resource.getTsFile().getName()).getTime();
-    } catch (final IOException e) {
-      LOGGER.warn(
-          "Pipe {}@{}: failed to get the generation time of TsFile {}, extract 
it anyway"
-              + " (historical data extraction time lower bound: {})",
-          pipeName,
-          dataRegionId,
-          resource.getTsFilePath(),
-          historicalDataExtractionTimeLowerBound,
-          e);
-      // If failed to get the generation time of the TsFile, we will extract 
the data in the TsFile
-      // anyway.
-      return true;
-    }
-  }
-
   @Override
   public synchronized Event supply() {
     if (!hasBeenStarted && 
StorageEngine.getInstance().isReadyForNonReadWriteFunctions()) {

Reply via email to