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()) {