diff --git a/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/event/dml/insertion/TsFileInsertionEvent.java b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/event/dml/insertion/TsFileInsertionEvent.java
index 4c7fffcfba517..594a3c89a648d 100644
--- a/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/event/dml/insertion/TsFileInsertionEvent.java
+++ b/iotdb-api/pipe-api/src/main/java/org/apache/iotdb/pipe/api/event/dml/insertion/TsFileInsertionEvent.java
@@ -33,6 +33,11 @@ public interface TsFileInsertionEvent extends Event, AutoCloseable {
* The method is used to convert the {@link TsFileInsertionEvent} into several {@link
* TabletInsertionEvent}s.
*
+ *
The returned iterable represents the lifetime of the conversion. Processors that retain it
+ * for asynchronous work must eventually exhaust the iterable before the event can be considered
+ * complete. A processor should not retain the iterable and silently abandon it, because the
+ * source cannot otherwise determine when all generated tablet events have finished.
+ *
* @return {@code Iterable} the list of {@link TabletInsertionEvent}
*/
Iterable toTabletInsertionEvents();
diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollector.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollector.java
index ac0f6c7bf0b74..53efadc77ed0d 100644
--- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollector.java
+++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollector.java
@@ -273,8 +273,18 @@ private void collectEvent(final Event event) {
((PipeHeartbeatEvent) event).recordConnectorQueueSize(pendingQueue);
}
+ if (event instanceof PipeRawTabletInsertionEvent
+ && ((PipeRawTabletInsertionEvent) event).getSourceEvent()
+ instanceof PipeTsFileInsertionEvent) {
+ // Register before publishing to the queue so a concurrent queue cleanup can report a
+ // generated tablet as discarded instead of racing with registration.
+ ((PipeTsFileInsertionEvent) ((PipeRawTabletInsertionEvent) event).getSourceEvent())
+ .registerGeneratedTabletInsertionEvent((PipeRawTabletInsertionEvent) event);
+ }
+
if (pendingQueue.offer(event)) {
collectInvocationCount.incrementAndGet();
+ return;
}
}
diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeRawTabletInsertionEvent.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeRawTabletInsertionEvent.java
index dc2ab1d381fd8..d113941a42cca 100644
--- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeRawTabletInsertionEvent.java
+++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tablet/PipeRawTabletInsertionEvent.java
@@ -64,6 +64,8 @@ public class PipeRawTabletInsertionEvent extends PipeInsertionEvent
private final boolean isAligned;
private final EnrichedEvent sourceEvent;
+ private final AtomicBoolean generatedEventRegisteredWithSource = new AtomicBoolean(false);
+ private final AtomicBoolean generatedEventOutcomeReported = new AtomicBoolean(false);
private boolean needToReport;
private final PipeTabletMemoryBlock allocatedMemoryBlock;
@@ -118,12 +120,19 @@ private PipeRawTabletInsertionEvent(
this.allocatedMemoryBlock =
PipeDataNodeResourceManager.memory().forceAllocateForTabletWithRetry(0);
- if (needToReport) {
+ if (needToReport || sourceEvent instanceof PipeTsFileInsertionEvent) {
addOnCommittedHook(
() -> {
- if (shouldReportOnCommit) {
+ // The last generated tablet still owns the progress-index cleanup. Other generated
+ // tablets only contribute to the transfer completion of their source TsFile.
+ if (shouldReportOnCommit && needToReport) {
eliminateProgressIndex();
}
+ if (sourceEvent instanceof PipeTsFileInsertionEvent
+ && generatedEventOutcomeReported.compareAndSet(false, true)) {
+ ((PipeTsFileInsertionEvent) sourceEvent)
+ .markGeneratedTabletInsertionEventAsTransferred();
+ }
});
}
}
@@ -306,9 +315,26 @@ public boolean internallyDecreaseResourceReferenceCount(final String holderMessa
return true;
}
+ /** Marks this generated tablet as accepted by the processor output collector. */
+ public void markAsGeneratedEventRegisteredWithSource() {
+ generatedEventRegisteredWithSource.set(true);
+ }
+
+ @Override
+ public boolean clearReferenceCount(final String holderMessage) {
+ final boolean cleared = super.clearReferenceCount(holderMessage);
+ if (generatedEventRegisteredWithSource.get()
+ && generatedEventOutcomeReported.compareAndSet(false, true)
+ && sourceEvent instanceof PipeTsFileInsertionEvent) {
+ ((PipeTsFileInsertionEvent) sourceEvent).markGeneratedTabletInsertionEventAsDiscarded();
+ }
+ return cleared;
+ }
+
protected void eliminateProgressIndex() {
if (sourceEvent instanceof PipeTsFileInsertionEvent) {
- ((PipeTsFileInsertionEvent) sourceEvent).eliminateProgressIndex();
+ final PipeTsFileInsertionEvent tsFileInsertionEvent = (PipeTsFileInsertionEvent) sourceEvent;
+ tsFileInsertionEvent.eliminateProgressIndex();
}
}
@@ -391,7 +417,7 @@ public boolean mayEventPathsOverlappedWithPattern() {
}
public void markAsNeedToReport() {
- if (!needToReport) {
+ if (!needToReport && !(sourceEvent instanceof PipeTsFileInsertionEvent)) {
addOnCommittedHook(
() -> {
if (shouldReportOnCommit) {
diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
index 335278e2f3790..883cea26bb355 100644
--- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
+++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/tsfile/PipeTsFileInsertionEvent.java
@@ -111,6 +111,18 @@ public class PipeTsFileInsertionEvent extends PipeInsertionEvent
private final AtomicInteger parsedTabletInsertionEventCount = new AtomicInteger(0);
private final AtomicBoolean isTsFileParsingCompleted = new AtomicBoolean(false);
private final AtomicLong parsedPointCountForCount = new AtomicLong(0);
+ private final AtomicInteger generatedTabletInsertionEventCount = new AtomicInteger(0);
+ private final AtomicInteger transferredGeneratedTabletInsertionEventCount = new AtomicInteger(0);
+ private final AtomicInteger discardedGeneratedTabletInsertionEventCount = new AtomicInteger(0);
+ private final AtomicBoolean generatedTabletInsertionEventsParsingStarted =
+ new AtomicBoolean(false);
+ private final AtomicBoolean generatedTabletInsertionEventsParsingCompleted =
+ new AtomicBoolean(false);
+ private final AtomicBoolean isTsFileEventCommitted = new AtomicBoolean(false);
+ private final AtomicBoolean isTransferred = new AtomicBoolean(false);
+ private final AtomicBoolean isDiscarded = new AtomicBoolean(false);
+ private final List onTransferredHooks = new ArrayList<>();
+ private final List onDiscardedHooks = new ArrayList<>();
// The point count of the TsFile. Used for metrics on IoTConsensusV2' receiver side.
// May be updated after it is flushed. Should be negative if not set.
@@ -298,6 +310,7 @@ private PipeTsFileInsertionEvent(
if (shouldReportOnCommit) {
eliminateProgressIndex();
}
+ markAsTransferredOnCommit();
});
}
@@ -477,6 +490,16 @@ public boolean internallyDecreaseResourceReferenceCount(final String holderMessa
}
}
+ @Override
+ public boolean clearReferenceCount(final String holderMessage) {
+ final boolean cleared = super.clearReferenceCount(holderMessage);
+ // clearReferenceCount intentionally bypasses the commit queue. Notify the realtime source even
+ // when the event was already released by a previous holder: its commit hook may still be
+ // waiting behind an earlier commit, while the explicit discard must release the in-flight slot.
+ markAsDiscarded();
+ return cleared;
+ }
+
@Override
public void bindProgressIndex(final ProgressIndex overridingProgressIndex) {
this.overridingProgressIndex = overridingProgressIndex;
@@ -532,6 +555,133 @@ public void eliminateProgressIndex() {
}
}
+ /**
+ * Records a tablet generated from this TsFile that has been accepted by the processor output
+ * collector. The transfer hook must wait for all such tablets, rather than the first one, before
+ * allowing the region-level downgrade to end.
+ */
+ public void registerGeneratedTabletInsertionEvent() {
+ generatedTabletInsertionEventCount.incrementAndGet();
+ }
+
+ public void registerGeneratedTabletInsertionEvent(final PipeRawTabletInsertionEvent event) {
+ event.markAsGeneratedEventRegisteredWithSource();
+ registerGeneratedTabletInsertionEvent();
+ }
+
+ /** Marks that a processor has obtained the tablet iterable and may consume it later. */
+ public void markGeneratedTabletInsertionEventsParsingStarted() {
+ generatedTabletInsertionEventsParsingStarted.set(true);
+ }
+
+ /** Records that one tablet generated from this TsFile has been committed downstream. */
+ public void markGeneratedTabletInsertionEventAsTransferred() {
+ transferredGeneratedTabletInsertionEventCount.incrementAndGet();
+ markAsTransferredIfGeneratedTabletEventsCompleted();
+ }
+
+ /** Records that one generated tablet was discarded before it could be committed downstream. */
+ public void markGeneratedTabletInsertionEventAsDiscarded() {
+ discardedGeneratedTabletInsertionEventCount.incrementAndGet();
+ markAsTransferredIfGeneratedTabletEventsCompleted();
+ }
+
+ /** Marks the end of tablet generation for this TsFile. */
+ public void markGeneratedTabletInsertionEventsParsingCompleted() {
+ generatedTabletInsertionEventsParsingCompleted.set(true);
+ markAsTransferredIfGeneratedTabletEventsCompleted();
+ }
+
+ private void markAsTransferredIfGeneratedTabletEventsCompleted() {
+ final boolean generationBoundaryReached =
+ !generatedTabletInsertionEventsParsingStarted.get()
+ || generatedTabletInsertionEventsParsingCompleted.get();
+ // The normal processor path does not retain a tablet iterable unless parsing has explicitly
+ // been marked as started. In that path, registering the generated tablets and reaching the
+ // generation boundary is sufficient to establish the transfer boundary (the source TsFile
+ // commit may be reported separately or may be intentionally skipped). Once an iterable has
+ // been retained, however, the source TsFile commit must also be observed before the generated
+ // tablet counts can release the TsFile. This distinction keeps deferred parsing safe without
+ // delaying the established non-deferred path.
+ if ((isTsFileEventCommitted.get() || !generatedTabletInsertionEventsParsingStarted.get())
+ && generationBoundaryReached
+ && transferredGeneratedTabletInsertionEventCount.get()
+ + discardedGeneratedTabletInsertionEventCount.get()
+ >= generatedTabletInsertionEventCount.get()) {
+ markAsTransferred();
+ }
+ }
+
+ private void markAsTransferredOnCommit() {
+ // A processor may retain the iterable returned by toTabletInsertionEvents() and consume it
+ // after process() returns. In that case the iterable wrapper keeps the generation boundary open
+ // until its iterator is exhausted. If no processor ever requests the iterable (for example,
+ // DoNothingProcessor), committing the TsFile itself is the completion boundary.
+ isTsFileEventCommitted.set(true);
+ markAsTransferredIfGeneratedTabletEventsCompleted();
+ }
+
+ private void markAsTransferred() {
+ final List hooksToRun;
+ synchronized (onTransferredHooks) {
+ if (isDiscarded.get() || !isTransferred.compareAndSet(false, true)) {
+ return;
+ }
+ hooksToRun = new ArrayList<>(onTransferredHooks);
+ onTransferredHooks.clear();
+ onDiscardedHooks.clear();
+ }
+ hooksToRun.forEach(Runnable::run);
+ }
+
+ private void markAsDiscarded() {
+ final List hooksToRun;
+ synchronized (onTransferredHooks) {
+ if (isTransferred.get() || !isDiscarded.compareAndSet(false, true)) {
+ return;
+ }
+ hooksToRun = new ArrayList<>(onDiscardedHooks);
+ onDiscardedHooks.clear();
+ onTransferredHooks.clear();
+ }
+ hooksToRun.forEach(Runnable::run);
+ }
+
+ /**
+ * Adds a hook that is invoked after this TsFile, or the last tablet generated from it, is
+ * committed downstream.
+ */
+ public void addOnTransferredHook(final Runnable hook) {
+ boolean runHook = false;
+ synchronized (onTransferredHooks) {
+ if (!isTransferred.get() && !isDiscarded.get()) {
+ onTransferredHooks.add(hook);
+ return;
+ }
+ runHook = isTransferred.get();
+ }
+ if (runHook) {
+ hook.run();
+ }
+ }
+
+ /**
+ * Adds a hook that runs when this TsFile event is explicitly discarded via clearReferenceCount.
+ */
+ public void addOnDiscardedHook(final Runnable hook) {
+ boolean runHook = false;
+ synchronized (onTransferredHooks) {
+ if (!isTransferred.get() && !isDiscarded.get()) {
+ onDiscardedHooks.add(hook);
+ return;
+ }
+ runHook = isDiscarded.get();
+ }
+ if (runHook) {
+ hook.run();
+ }
+ }
+
public PipeTsFileInsertionEvent bindTsFileDedupScopeID(final String tsFileDedupScopeID) {
this.tsFileDedupScopeID = tsFileDedupScopeID;
return this;
@@ -802,6 +952,7 @@ public void consumeTabletInsertionEventsWithRetry(
final String callerName,
final PipeProcessorSubtaskExecutionGuard processorExecutionGuard)
throws Exception {
+ markGeneratedTabletInsertionEventsParsingStarted();
try {
while (true) {
processorExecutionGuard.check();
@@ -809,6 +960,7 @@ public void consumeTabletInsertionEventsWithRetry(
getNextTabletInsertionEventFromSavedProgress(processorExecutionGuard);
if (parsedEvent == null) {
isTsFileParsingCompleted.set(true);
+ markGeneratedTabletInsertionEventsParsingCompleted();
releaseTsFileParserMemoryIfReserved();
return;
}
@@ -981,14 +1133,18 @@ public Iterable toTabletInsertionEvents() throws PipeExcep
public Iterable toTabletInsertionEvents(final long timeoutMs)
throws PipeException {
+ markGeneratedTabletInsertionEventsParsingStarted();
+
try {
if (!waitForTsFileClose()) {
LOGGER.warn(DataNodePipeMessages.PIPE_SKIPPING_TEMPORARY_TSFILE_S_PARSING_WHICH, tsFile);
+ markGeneratedTabletInsertionEventsParsingCompleted();
return Collections.emptyList();
}
waitForResourceEnough4Parsing(timeoutMs);
- return initEventParser().toTabletInsertionEvents();
+ return wrapGeneratedTabletInsertionEvents(initEventParser().toTabletInsertionEvents());
} catch (final Exception e) {
+ markGeneratedTabletInsertionEventsParsingCompleted();
close();
// close() should be called before re-interrupting the thread
@@ -1014,6 +1170,49 @@ public Iterable toTabletInsertionEvents(final long timeout
}
}
+ private Iterable wrapGeneratedTabletInsertionEvents(
+ final Iterable iterable) {
+ return () -> {
+ final Iterator iterator;
+ try {
+ iterator = iterable.iterator();
+ } catch (final RuntimeException e) {
+ markGeneratedTabletInsertionEventsParsingCompleted();
+ throw e;
+ }
+ return new Iterator() {
+ @Override
+ public boolean hasNext() {
+ try {
+ final boolean hasNext = iterator.hasNext();
+ if (!hasNext) {
+ markGeneratedTabletInsertionEventsParsingCompleted();
+ }
+ return hasNext;
+ } catch (final RuntimeException e) {
+ markGeneratedTabletInsertionEventsParsingCompleted();
+ throw e;
+ }
+ }
+
+ @Override
+ public TabletInsertionEvent next() {
+ try {
+ return iterator.next();
+ } catch (final RuntimeException e) {
+ markGeneratedTabletInsertionEventsParsingCompleted();
+ throw e;
+ }
+ }
+
+ @Override
+ public void remove() {
+ iterator.remove();
+ }
+ };
+ };
+ }
+
private void reserveResource4Parsing(
final PipeProcessorSubtaskExecutionGuard processorExecutionGuard)
throws InterruptedException {
diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionHybridSource.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionHybridSource.java
index 65d03919ab1bc..6e56e6904f710 100644
--- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionHybridSource.java
+++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionHybridSource.java
@@ -21,6 +21,7 @@
import org.apache.iotdb.commons.exception.pipe.PipeRuntimeNonCriticalException;
import org.apache.iotdb.commons.pipe.config.PipeConfig;
+import org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant;
import org.apache.iotdb.commons.pipe.event.ProgressReportEvent;
import org.apache.iotdb.db.i18n.DataNodePipeMessages;
import org.apache.iotdb.db.pipe.agent.PipeDataNodeAgent;
@@ -33,6 +34,9 @@
import org.apache.iotdb.db.pipe.resource.PipeDataNodeResourceManager;
import org.apache.iotdb.db.pipe.source.dataregion.realtime.assigner.PipeTsFileEpochProgressIndexKeeper;
import org.apache.iotdb.db.pipe.source.dataregion.realtime.epoch.TsFileEpoch;
+import org.apache.iotdb.pipe.api.customizer.configuration.PipeExtractorRuntimeConfiguration;
+import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameterValidator;
+import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters;
import org.apache.iotdb.pipe.api.event.Event;
import org.apache.iotdb.pipe.api.event.dml.insertion.TabletInsertionEvent;
import org.apache.iotdb.pipe.api.event.dml.insertion.TsFileInsertionEvent;
@@ -40,7 +44,10 @@
import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
+import java.util.ArrayDeque;
+import java.util.Arrays;
import java.util.Collections;
+import java.util.Deque;
import java.util.Optional;
import java.util.Set;
import java.util.concurrent.ConcurrentHashMap;
@@ -50,13 +57,68 @@ public class PipeRealtimeDataRegionHybridSource extends PipeRealtimeDataRegionSo
private static final Logger LOGGER =
LoggerFactory.getLogger(PipeRealtimeDataRegionHybridSource.class);
+ private boolean isRegionLevelDowngradingEnabled =
+ PipeSourceConstant.EXTRACTOR_REALTIME_REGION_LEVEL_DOWNGRADING_DEFAULT_VALUE;
+ private final Object regionLevelDowngradingLock = new Object();
+
private final Set activeTsFileEpochs =
Collections.newSetFromMap(new ConcurrentHashMap<>());
private final Set degradedTsFileEpochs =
Collections.newSetFromMap(new ConcurrentHashMap<>());
+ private final Deque regionLevelDegradedTsFileEpochs = new ArrayDeque<>();
+ private final Deque eventsBeforeRegionLevelDowngrading = new ArrayDeque<>();
+ private final Deque regionLevelBufferedEvents = new ArrayDeque<>();
+
+ private volatile boolean isRegionLevelDegraded = false;
+ private TsFileEpoch regionLevelTailTsFileEpoch = null;
+ private boolean canSupplyEventsBeforeRegionLevelDowngrading = false;
+ private int inFlightTsFileCount = 0;
+
+ @Override
+ public void validate(final PipeParameterValidator validator) throws Exception {
+ super.validate(validator);
+ validator
+ .validateAttributeValueRange(
+ PipeSourceConstant.EXTRACTOR_REALTIME_REGION_LEVEL_DOWNGRADING_KEY,
+ true,
+ Boolean.TRUE.toString(),
+ Boolean.FALSE.toString())
+ .validateAttributeValueRange(
+ PipeSourceConstant.SOURCE_REALTIME_REGION_LEVEL_DOWNGRADING_KEY,
+ true,
+ Boolean.TRUE.toString(),
+ Boolean.FALSE.toString());
+ }
+
+ @Override
+ public void customize(
+ final PipeParameters parameters, final PipeExtractorRuntimeConfiguration configuration)
+ throws Exception {
+ super.customize(parameters, configuration);
+ isRegionLevelDowngradingEnabled =
+ parameters.getBooleanOrDefault(
+ Arrays.asList(
+ PipeSourceConstant.EXTRACTOR_REALTIME_REGION_LEVEL_DOWNGRADING_KEY,
+ PipeSourceConstant.SOURCE_REALTIME_REGION_LEVEL_DOWNGRADING_KEY),
+ PipeSourceConstant.EXTRACTOR_REALTIME_REGION_LEVEL_DOWNGRADING_DEFAULT_VALUE);
+ }
@Override
protected void doExtract(final PipeRealtimeEvent event) {
+ if (isRegionLevelDowngradingEnabled) {
+ synchronized (regionLevelDowngradingLock) {
+ if (isClosed.get()) {
+ event.decreaseReferenceCount(PipeRealtimeDataRegionHybridSource.class.getName(), false);
+ return;
+ }
+ doExtractInternal(event);
+ }
+ return;
+ }
+ doExtractInternal(event);
+ }
+
+ private void doExtractInternal(final PipeRealtimeEvent event) {
final Event eventToExtract = event.getEvent();
if (eventToExtract instanceof TabletInsertionEvent) {
@@ -91,9 +153,46 @@ public boolean isNeedListenToInsertNode() {
private void extractTabletInsertion(final PipeRealtimeEvent event) {
markTsFileEpochActive(event.getTsFileEpoch());
+ if (isRegionLevelDowngradingEnabled
+ && isRegionLevelDegraded
+ && degradedTsFileEpochs.contains(event.getTsFileEpoch())) {
+ event.getTsFileEpoch().migrateState(this, currentState -> TsFileEpoch.State.USING_TSFILE);
+ PipeTsFileEpochProgressIndexKeeper.getInstance()
+ .registerProgressIndex(
+ dataRegionId, getTsFileDedupScopeID(), event.getTsFileEpoch().getResource());
+ event.decreaseReferenceCount(PipeRealtimeDataRegionHybridSource.class.getName(), false);
+ return;
+ }
+
TsFileEpoch.State state;
- if (canNotUseTabletAnymore(event)) {
+ if (isRegionLevelDowngradingEnabled && isRegionLevelDegraded) {
+ // Retain only the newest epoch as a realtime tail. Once another epoch arrives, the previous
+ // tail is downgraded to bound the buffered tablet memory to roughly one TsFile.
+ prepareRegionLevelTailTsFileEpochUnderLock(event.getTsFileEpoch());
+ if (canNotUseTabletAnymore(event)) {
+ if (regionLevelTailTsFileEpoch == event.getTsFileEpoch()) {
+ promoteRegionLevelTailTsFileEpochUnderLock();
+ bufferPendingEventsForRegionLevelExitUnderLock();
+ rebalanceRegionLevelBufferedEventsUnderLock();
+ }
+ event.getTsFileEpoch().migrateState(this, currentState -> TsFileEpoch.State.USING_TSFILE);
+ PipeTsFileEpochProgressIndexKeeper.getInstance()
+ .registerProgressIndex(
+ dataRegionId, getTsFileDedupScopeID(), event.getTsFileEpoch().getResource());
+ markTsFileEpochDegradedFromExtraction(event.getTsFileEpoch());
+ event.decreaseReferenceCount(PipeRealtimeDataRegionHybridSource.class.getName(), false);
+ return;
+ }
+ event
+ .getTsFileEpoch()
+ .migrateState(
+ this,
+ currentState ->
+ currentState == TsFileEpoch.State.EMPTY
+ ? TsFileEpoch.State.USING_TABLET
+ : currentState);
+ } else if (canNotUseTabletAnymore(event)) {
event.getTsFileEpoch().migrateState(this, curState -> TsFileEpoch.State.USING_TSFILE);
PipeTsFileEpochProgressIndexKeeper.getInstance()
.registerProgressIndex(
@@ -118,21 +217,25 @@ private void extractTabletInsertion(final PipeRealtimeEvent event) {
state = event.getTsFileEpoch().getState(this);
if (state == TsFileEpoch.State.USING_TSFILE || state == TsFileEpoch.State.USING_BOTH) {
- markTsFileEpochDegraded(event.getTsFileEpoch());
+ markTsFileEpochDegradedFromExtraction(event.getTsFileEpoch());
}
switch (state) {
case USING_TSFILE:
// Ignore the tablet event.
event.decreaseReferenceCount(PipeRealtimeDataRegionHybridSource.class.getName(), false);
break;
- case EMPTY:
- case USING_TABLET:
case USING_BOTH:
- // USING_BOTH indicates that there are discarded events previously.
- // In this case, we need to delay the progress report to tsFile event, to avoid losing data.
- if (state == TsFileEpoch.State.USING_BOTH) {
- event.skipReportOnCommit();
+ if (isRegionLevelDowngradingEnabled && isRegionLevelDegraded) {
+ event.decreaseReferenceCount(PipeRealtimeDataRegionHybridSource.class.getName(), false);
+ break;
}
+ // USING_BOTH indicates that there are discarded events previously. In this case, we need
+ // to delay the progress report to the TsFile event, to avoid losing data.
+ event.skipReportOnCommit();
+ pendingQueue.offer(event);
+ break;
+ case EMPTY:
+ case USING_TABLET:
pendingQueue.offer(event);
break;
default:
@@ -173,11 +276,18 @@ private void extractTsFileInsertion(final PipeRealtimeEvent event) {
});
final TsFileEpoch.State state = event.getTsFileEpoch().getState(this);
- if (state == TsFileEpoch.State.USING_BOTH) {
- markTsFileEpochDegraded(event.getTsFileEpoch());
+ if (state == TsFileEpoch.State.USING_BOTH
+ || (isRegionLevelDowngradingEnabled
+ && isRegionLevelDegraded
+ && state == TsFileEpoch.State.USING_TSFILE)) {
+ markTsFileEpochDegradedFromExtraction(event.getTsFileEpoch());
}
switch (state) {
case USING_TABLET:
+ if (isRegionLevelDowngradingEnabled && isRegionLevelDegraded) {
+ pendingQueue.offer(event);
+ return;
+ }
// If the state is USING_TABLET, discard the event
PipeTsFileEpochProgressIndexKeeper.getInstance()
.eliminateProgressIndex(
@@ -201,43 +311,237 @@ private void extractTsFileInsertion(final PipeRealtimeEvent event) {
}
private void markTsFileEpochActive(final TsFileEpoch tsFileEpoch) {
- activeTsFileEpochs.add(tsFileEpoch);
- reportTsFileEpochDegradedStatus();
+ synchronized (regionLevelDowngradingLock) {
+ activeTsFileEpochs.add(tsFileEpoch);
+ reportTsFileEpochDegradedStatusUnderLock();
+ }
}
private void markTsFileEpochDegraded(final TsFileEpoch tsFileEpoch) {
+ markTsFileEpochDegraded(tsFileEpoch, false);
+ }
+
+ private void markTsFileEpochDegradedFromExtraction(final TsFileEpoch tsFileEpoch) {
+ markTsFileEpochDegraded(tsFileEpoch, true);
+ }
+
+ private void markTsFileEpochDegraded(
+ final TsFileEpoch tsFileEpoch, final boolean shouldPreservePendingEvents) {
+ synchronized (regionLevelDowngradingLock) {
+ final boolean wasRegionLevelDegraded = isRegionLevelDegraded;
+ if (isRegionLevelDowngradingEnabled
+ && shouldPreservePendingEvents
+ && !wasRegionLevelDegraded) {
+ PipeRealtimeEvent pendingEvent;
+ while ((pendingEvent = (PipeRealtimeEvent) pendingQueue.directPoll()) != null) {
+ eventsBeforeRegionLevelDowngrading.offerLast(pendingEvent);
+ }
+ canSupplyEventsBeforeRegionLevelDowngrading = true;
+ }
+ if (regionLevelTailTsFileEpoch == tsFileEpoch) {
+ regionLevelTailTsFileEpoch = null;
+ }
+ markTsFileEpochDegradedUnderLock(tsFileEpoch);
+ if (isRegionLevelDowngradingEnabled) {
+ // A downgrade discovered while supplying an event happens after all remaining pending
+ // events. Track those events as the possible realtime tail before rebalancing the queues.
+ if (!wasRegionLevelDegraded && !shouldPreservePendingEvents) {
+ bufferPendingEventsAndTrackRegionLevelTailUnderLock();
+ }
+ rebalanceRegionLevelBufferedEventsUnderLock();
+ }
+ }
+ }
+
+ private void markTsFileEpochDegradedUnderLock(final TsFileEpoch tsFileEpoch) {
activeTsFileEpochs.add(tsFileEpoch);
- degradedTsFileEpochs.add(tsFileEpoch);
- reportTsFileEpochDegradedStatus();
+ if (degradedTsFileEpochs.add(tsFileEpoch) && isRegionLevelDowngradingEnabled) {
+ regionLevelDegradedTsFileEpochs.offerLast(tsFileEpoch);
+ }
+ if (isRegionLevelDowngradingEnabled) {
+ isRegionLevelDegraded = true;
+ }
+ reportTsFileEpochDegradedStatusUnderLock();
+ }
+
+ private void prepareRegionLevelTailTsFileEpochUnderLock(final TsFileEpoch tsFileEpoch) {
+ if (regionLevelTailTsFileEpoch == tsFileEpoch) {
+ return;
+ }
+
+ if (regionLevelTailTsFileEpoch != null) {
+ promoteRegionLevelTailTsFileEpochUnderLock();
+ bufferPendingEventsForRegionLevelExitUnderLock();
+ rebalanceRegionLevelBufferedEventsUnderLock();
+ }
+ regionLevelTailTsFileEpoch = tsFileEpoch;
+ }
+
+ private void promoteRegionLevelTailTsFileEpochUnderLock() {
+ final TsFileEpoch tsFileEpoch = regionLevelTailTsFileEpoch;
+ if (tsFileEpoch == null) {
+ return;
+ }
+
+ regionLevelTailTsFileEpoch = null;
+ tsFileEpoch.migrateState(this, state -> TsFileEpoch.State.USING_TSFILE);
+ PipeTsFileEpochProgressIndexKeeper.getInstance()
+ .registerProgressIndex(dataRegionId, getTsFileDedupScopeID(), tsFileEpoch.getResource());
+ markTsFileEpochDegradedUnderLock(tsFileEpoch);
+ }
+
+ private void bufferPendingEventsAndTrackRegionLevelTailUnderLock() {
+ PipeRealtimeEvent pendingEvent;
+ while ((pendingEvent = (PipeRealtimeEvent) pendingQueue.directPoll()) != null) {
+ if (pendingEvent.getEvent() instanceof TabletInsertionEvent
+ && !degradedTsFileEpochs.contains(pendingEvent.getTsFileEpoch())) {
+ if (regionLevelTailTsFileEpoch != null
+ && regionLevelTailTsFileEpoch != pendingEvent.getTsFileEpoch()) {
+ promoteRegionLevelTailTsFileEpochUnderLock();
+ }
+ regionLevelTailTsFileEpoch = pendingEvent.getTsFileEpoch();
+ }
+ regionLevelBufferedEvents.offerLast(pendingEvent);
+ }
+ }
+
+ private void rebalanceRegionLevelBufferedEventsUnderLock() {
+ bufferPendingEventsForRegionLevelExitUnderLock();
+
+ final TsFileEpoch nextDegradedTsFileEpoch = regionLevelDegradedTsFileEpochs.peekFirst();
+ final Deque retainedEvents = new ArrayDeque<>();
+ boolean nextDegradedTsFileEventPromoted = false;
+ PipeRealtimeEvent bufferedEvent;
+ while ((bufferedEvent = regionLevelBufferedEvents.pollFirst()) != null) {
+ if (bufferedEvent.getEvent() instanceof TabletInsertionEvent
+ && degradedTsFileEpochs.contains(bufferedEvent.getTsFileEpoch())) {
+ bufferedEvent.decreaseReferenceCount(
+ PipeRealtimeDataRegionHybridSource.class.getName(), false);
+ } else if (!nextDegradedTsFileEventPromoted
+ && bufferedEvent.getEvent() instanceof TsFileInsertionEvent
+ && bufferedEvent.getTsFileEpoch() == nextDegradedTsFileEpoch) {
+ pendingQueue.offer(bufferedEvent);
+ nextDegradedTsFileEventPromoted = true;
+ } else {
+ retainedEvents.offerLast(bufferedEvent);
+ }
+ }
+ regionLevelBufferedEvents.addAll(retainedEvents);
}
private void clearTsFileEpoch(final TsFileEpoch tsFileEpoch) {
+ synchronized (regionLevelDowngradingLock) {
+ clearTsFileEpochUnderLock(tsFileEpoch);
+ }
+ }
+
+ private void clearTsFileEpochAfterCommit(final TsFileEpoch tsFileEpoch) {
+ synchronized (regionLevelDowngradingLock) {
+ if (isClosed.get()) {
+ return;
+ }
+ if (inFlightTsFileCount > 0) {
+ --inFlightTsFileCount;
+ }
+ clearTsFileEpochUnderLock(tsFileEpoch);
+
+ // The decision whether newer writes can resume the realtime path must be made at the exact
+ // point when the last currently degraded TsFile is committed. Otherwise, writes arriving
+ // between the commit and the next supply call would still be unnecessarily downgraded.
+ if (isRegionLevelDegraded
+ && inFlightTsFileCount == 0
+ && degradedTsFileEpochs.isEmpty()
+ && eventsBeforeRegionLevelDowngrading.isEmpty()) {
+ bufferPendingEventsForRegionLevelExitUnderLock();
+ tryExitRegionLevelDowngrading(false);
+ }
+ }
+ }
+
+ private void bufferPendingEventsForRegionLevelExitUnderLock() {
+ PipeRealtimeEvent event;
+ while ((event = (PipeRealtimeEvent) pendingQueue.directPoll()) != null) {
+ regionLevelBufferedEvents.offerLast(event);
+ }
+ }
+
+ private void clearTsFileEpochUnderLock(final TsFileEpoch tsFileEpoch) {
activeTsFileEpochs.remove(tsFileEpoch);
degradedTsFileEpochs.remove(tsFileEpoch);
- reportTsFileEpochDegradedStatus();
+ regionLevelDegradedTsFileEpochs.remove(tsFileEpoch);
+ if (regionLevelTailTsFileEpoch == tsFileEpoch) {
+ regionLevelTailTsFileEpoch = null;
+ }
+ if (isRegionLevelDowngradingEnabled && isRegionLevelDegraded && inFlightTsFileCount == 0) {
+ rebalanceRegionLevelBufferedEventsUnderLock();
+ }
+ reportTsFileEpochDegradedStatusUnderLock();
+ }
+
+ private void clearRegionLevelBufferedEventsUnderLock() {
+ clearBufferedEventsUnderLock(eventsBeforeRegionLevelDowngrading);
+ clearBufferedEventsUnderLock(regionLevelBufferedEvents);
+ }
+
+ private void clearBufferedEventsUnderLock(final Deque events) {
+ PipeRealtimeEvent event;
+ while ((event = events.pollFirst()) != null) {
+ event.clearReferenceCount(PipeRealtimeDataRegionHybridSource.class.getName());
+ }
}
- private void reportTsFileEpochDegradedStatus() {
- if (activeTsFileEpochs.isEmpty()) {
+ private void reportTsFileEpochDegradedStatusUnderLock() {
+ if (isRegionLevelDowngradingEnabled && isRegionLevelDegraded) {
+ PipeDataNodeAgent.task()
+ .setPipeTsFileEpochDegraded(pipeName, creationTime, dataRegionId, true);
+ } else if (activeTsFileEpochs.isEmpty()) {
PipeDataNodeAgent.task().clearPipeTsFileEpochDegraded(pipeName, creationTime, dataRegionId);
} else {
PipeDataNodeAgent.task()
.setPipeTsFileEpochDegraded(
- pipeName, creationTime, dataRegionId, !degradedTsFileEpochs.isEmpty());
+ pipeName,
+ creationTime,
+ dataRegionId,
+ isRegionLevelDowngradingEnabled ? false : !degradedTsFileEpochs.isEmpty());
}
}
@Override
public void close() throws Exception {
try {
+ // Do not hold regionLevelDowngradingLock while waiting for the assigner to stop. An event
+ // already being assigned may need the same lock to finish extraction.
super.close();
} finally {
- activeTsFileEpochs.clear();
- degradedTsFileEpochs.clear();
- PipeDataNodeAgent.task().clearPipeTsFileEpochDegraded(pipeName, creationTime, dataRegionId);
+ synchronized (regionLevelDowngradingLock) {
+ clearRegionLevelBufferedEventsUnderLock();
+ activeTsFileEpochs.clear();
+ degradedTsFileEpochs.clear();
+ regionLevelDegradedTsFileEpochs.clear();
+ isRegionLevelDegraded = false;
+ regionLevelTailTsFileEpoch = null;
+ canSupplyEventsBeforeRegionLevelDowngrading = false;
+ inFlightTsFileCount = 0;
+ PipeDataNodeAgent.task().clearPipeTsFileEpochDegraded(pipeName, creationTime, dataRegionId);
+ }
}
}
+ @Override
+ protected void extractProgressReportEvent(final PipeRealtimeEvent event) {
+ if (isRegionLevelDowngradingEnabled) {
+ synchronized (regionLevelDowngradingLock) {
+ if (isClosed.get()) {
+ event.decreaseReferenceCount(PipeRealtimeDataRegionHybridSource.class.getName(), false);
+ return;
+ }
+ super.extractProgressReportEvent(event);
+ }
+ return;
+ }
+ super.extractProgressReportEvent(event);
+ }
+
// If the insertNode's memory has reached the dangerous threshold, we should not extract any
// tablets.
private boolean canNotUseTabletAnymore(final PipeRealtimeEvent event) {
@@ -278,47 +582,199 @@ private boolean canNotUseTabletAnymore(final PipeRealtimeEvent event) {
@Override
public Event supply() {
+ if (isRegionLevelDowngradingEnabled) {
+ synchronized (regionLevelDowngradingLock) {
+ return isRegionLevelDegraded ? supplyRegionLevelDegradedInternal() : supplyInternal();
+ }
+ }
+ return supplyInternal();
+ }
+
+ private Event supplyInternal() {
PipeRealtimeEvent realtimeEvent = (PipeRealtimeEvent) pendingQueue.directPoll();
while (realtimeEvent != null) {
- Event suppliedEvent;
+ final Event suppliedEvent = supplyExtractedEvent(realtimeEvent);
+ if (suppliedEvent != null) {
+ return suppliedEvent;
+ }
+
+ if (isRegionLevelDowngradingEnabled && isRegionLevelDegraded) {
+ return supplyRegionLevelDegradedInternal();
+ }
+
+ realtimeEvent = (PipeRealtimeEvent) pendingQueue.directPoll();
+ }
+
+ // Means the pending queue is empty.
+ return null;
+ }
+
+ private Event supplyExtractedEvent(final PipeRealtimeEvent realtimeEvent) {
+ Event suppliedEvent;
+
+ // Used to judge the type of the event, not directly for supplying.
+ final Event eventToSupply = realtimeEvent.getEvent();
+ if (eventToSupply instanceof TabletInsertionEvent) {
+ suppliedEvent = supplyTabletInsertion(realtimeEvent);
+ } else if (eventToSupply instanceof TsFileInsertionEvent) {
+ suppliedEvent = supplyTsFileInsertion(realtimeEvent);
+ } else if (eventToSupply instanceof PipeHeartbeatEvent) {
+ suppliedEvent = supplyHeartbeat(realtimeEvent);
+ } else if (eventToSupply instanceof PipeDeleteDataNodeEvent
+ || eventToSupply instanceof ProgressReportEvent) {
+ suppliedEvent = supplyDirectly(realtimeEvent);
+ } else {
+ throw new UnsupportedOperationException(
+ String.format(
+ DataNodePipeMessages
+ .PIPE_EXCEPTION_UNSUPPORTED_EVENT_TYPE_S_FOR_HYBRID_REALTIME_EXTRACTOR_S_474BAAC2,
+ eventToSupply.getClass(),
+ this));
+ }
+
+ realtimeEvent.decreaseReferenceCount(PipeRealtimeDataRegionHybridSource.class.getName(), false);
+
+ if (suppliedEvent != null) {
+ suppliedEvent = assignReplicateIndexIfNeeded(realtimeEvent, suppliedEvent);
+ maySkipIndex4Event(realtimeEvent);
+ }
+ return suppliedEvent;
+ }
+
+ private Event supplyRegionLevelDegradedInternal() {
+ // Once region-level downgrading starts, only one TsFile is allowed to be in flight. Events of
+ // newer epochs are buffered until all currently degraded TsFiles are committed downstream.
+ if (inFlightTsFileCount > 0) {
+ return null;
+ }
+
+ final Event eventBeforeDowngrading = supplyEventsBeforeRegionLevelDowngradingInternal();
+ if (eventBeforeDowngrading != null) {
+ return eventBeforeDowngrading;
+ }
- // Used to judge the type of the event, not directly for supplying.
+ PipeRealtimeEvent realtimeEvent = (PipeRealtimeEvent) pendingQueue.directPoll();
+ while (realtimeEvent != null) {
final Event eventToSupply = realtimeEvent.getEvent();
if (eventToSupply instanceof TabletInsertionEvent) {
- suppliedEvent = supplyTabletInsertion(realtimeEvent);
+ if (degradedTsFileEpochs.contains(realtimeEvent.getTsFileEpoch())) {
+ realtimeEvent.decreaseReferenceCount(
+ PipeRealtimeDataRegionHybridSource.class.getName(), false);
+ } else {
+ regionLevelBufferedEvents.offerLast(realtimeEvent);
+ }
} else if (eventToSupply instanceof TsFileInsertionEvent) {
- suppliedEvent = supplyTsFileInsertion(realtimeEvent);
- } else if (eventToSupply instanceof PipeHeartbeatEvent) {
- suppliedEvent = supplyHeartbeat(realtimeEvent);
- } else if (eventToSupply instanceof PipeDeleteDataNodeEvent
- || eventToSupply instanceof ProgressReportEvent) {
- suppliedEvent = supplyDirectly(realtimeEvent);
+ final TsFileEpoch.State state = realtimeEvent.getTsFileEpoch().getState(this);
+ if (degradedTsFileEpochs.contains(realtimeEvent.getTsFileEpoch())
+ || state == TsFileEpoch.State.USING_TSFILE
+ || state == TsFileEpoch.State.USING_BOTH) {
+ markTsFileEpochDegraded(realtimeEvent.getTsFileEpoch());
+ if (regionLevelDegradedTsFileEpochs.peekFirst() != realtimeEvent.getTsFileEpoch()) {
+ regionLevelBufferedEvents.offerLast(realtimeEvent);
+ realtimeEvent = (PipeRealtimeEvent) pendingQueue.directPoll();
+ continue;
+ }
+ final Event suppliedEvent = supplyExtractedEvent(realtimeEvent);
+ if (suppliedEvent != null) {
+ return suppliedEvent;
+ }
+ } else {
+ regionLevelBufferedEvents.offerLast(realtimeEvent);
+ }
} else {
- throw new UnsupportedOperationException(
- String.format(
- DataNodePipeMessages
- .PIPE_EXCEPTION_UNSUPPORTED_EVENT_TYPE_S_FOR_HYBRID_REALTIME_EXTRACTOR_S_474BAAC2,
- eventToSupply.getClass(),
- this));
+ regionLevelBufferedEvents.offerLast(realtimeEvent);
}
- realtimeEvent.decreaseReferenceCount(
- PipeRealtimeDataRegionHybridSource.class.getName(), false);
+ realtimeEvent = (PipeRealtimeEvent) pendingQueue.directPoll();
+ }
- if (suppliedEvent != null) {
- suppliedEvent = assignReplicateIndexIfNeeded(realtimeEvent, suppliedEvent);
- maySkipIndex4Event(realtimeEvent);
- return suppliedEvent;
+ return tryExitRegionLevelDowngrading(true);
+ }
+
+ private Event supplyEventsBeforeRegionLevelDowngradingInternal() {
+ PipeRealtimeEvent realtimeEvent;
+ while ((realtimeEvent = eventsBeforeRegionLevelDowngrading.pollFirst()) != null) {
+ final Event eventToSupply = realtimeEvent.getEvent();
+
+ if (canSupplyEventsBeforeRegionLevelDowngrading) {
+ if (eventToSupply instanceof TabletInsertionEvent
+ && degradedTsFileEpochs.contains(realtimeEvent.getTsFileEpoch())) {
+ realtimeEvent.decreaseReferenceCount(
+ PipeRealtimeDataRegionHybridSource.class.getName(), false);
+ continue;
+ }
+
+ final Event suppliedEvent = supplyExtractedEvent(realtimeEvent);
+ if (suppliedEvent != null) {
+ return suppliedEvent;
+ }
+ if (eventToSupply instanceof TabletInsertionEvent
+ && degradedTsFileEpochs.contains(realtimeEvent.getTsFileEpoch())) {
+ canSupplyEventsBeforeRegionLevelDowngrading = false;
+ }
+ continue;
}
- realtimeEvent = (PipeRealtimeEvent) pendingQueue.directPoll();
+ if (eventToSupply instanceof TabletInsertionEvent) {
+ if (degradedTsFileEpochs.contains(realtimeEvent.getTsFileEpoch())) {
+ realtimeEvent.decreaseReferenceCount(
+ PipeRealtimeDataRegionHybridSource.class.getName(), false);
+ } else {
+ regionLevelBufferedEvents.offerLast(realtimeEvent);
+ }
+ } else if (eventToSupply instanceof TsFileInsertionEvent) {
+ final TsFileEpoch.State state = realtimeEvent.getTsFileEpoch().getState(this);
+ if (degradedTsFileEpochs.contains(realtimeEvent.getTsFileEpoch())
+ || state == TsFileEpoch.State.USING_TSFILE
+ || state == TsFileEpoch.State.USING_BOTH) {
+ markTsFileEpochDegraded(realtimeEvent.getTsFileEpoch());
+ if (regionLevelDegradedTsFileEpochs.peekFirst() != realtimeEvent.getTsFileEpoch()) {
+ regionLevelBufferedEvents.offerLast(realtimeEvent);
+ continue;
+ }
+ final Event suppliedEvent = supplyExtractedEvent(realtimeEvent);
+ if (suppliedEvent != null) {
+ return suppliedEvent;
+ }
+ } else {
+ regionLevelBufferedEvents.offerLast(realtimeEvent);
+ }
+ } else {
+ regionLevelBufferedEvents.offerLast(realtimeEvent);
+ }
}
- // Means the pending queue is empty.
+ canSupplyEventsBeforeRegionLevelDowngrading = false;
return null;
}
+ private Event tryExitRegionLevelDowngrading(final boolean shouldSupplyAfterTransition) {
+ if (!degradedTsFileEpochs.isEmpty()) {
+ // Some degraded epochs are still waiting for their TsFile to be flushed.
+ return null;
+ }
+
+ if (regionLevelTailTsFileEpoch != null
+ && regionLevelBufferedEvents.stream()
+ .filter(event -> event.getTsFileEpoch() == regionLevelTailTsFileEpoch)
+ .filter(event -> event.getEvent() instanceof TabletInsertionEvent)
+ .anyMatch(event -> event.getEvent().isReleased())) {
+ promoteRegionLevelTailTsFileEpochUnderLock();
+ rebalanceRegionLevelBufferedEventsUnderLock();
+ return shouldSupplyAfterTransition ? supplyRegionLevelDegradedInternal() : null;
+ }
+
+ isRegionLevelDegraded = false;
+ regionLevelTailTsFileEpoch = null;
+ PipeRealtimeEvent bufferedEvent;
+ while ((bufferedEvent = regionLevelBufferedEvents.pollFirst()) != null) {
+ pendingQueue.offer(bufferedEvent);
+ }
+ reportTsFileEpochDegradedStatusUnderLock();
+ return shouldSupplyAfterTransition ? supplyInternal() : null;
+ }
+
private Event supplyTabletInsertion(final PipeRealtimeEvent event) {
if (event.increaseReferenceCount(PipeRealtimeDataRegionHybridSource.class.getName())) {
return event.getEvent();
@@ -326,7 +782,14 @@ private Event supplyTabletInsertion(final PipeRealtimeEvent event) {
// If the event's reference count can not be increased, it means the data represented by
// this event is not reliable anymore. but the data represented by this event
// has been carried by the following tsfile event, so we can just discard this event.
- event.getTsFileEpoch().migrateState(this, s -> TsFileEpoch.State.USING_BOTH);
+ event
+ .getTsFileEpoch()
+ .migrateState(
+ this,
+ state ->
+ isRegionLevelDowngradingEnabled
+ ? TsFileEpoch.State.USING_TSFILE
+ : TsFileEpoch.State.USING_BOTH);
markTsFileEpochDegraded(event.getTsFileEpoch());
LOGGER.warn(DataNodePipeMessages.DISCARD_TABLET_EVENT_BECAUSE_IT_IS_NOT, event);
return null;
@@ -334,27 +797,44 @@ private Event supplyTabletInsertion(final PipeRealtimeEvent event) {
}
private Event supplyTsFileInsertion(final PipeRealtimeEvent event) {
- try {
- if (event.increaseReferenceCount(PipeRealtimeDataRegionHybridSource.class.getName())) {
- return event.getEvent();
+ if (isRegionLevelDowngradingEnabled
+ && event.getTsFileEpoch().getState(this) == TsFileEpoch.State.USING_TABLET) {
+ PipeTsFileEpochProgressIndexKeeper.getInstance()
+ .eliminateProgressIndex(
+ dataRegionId, getTsFileDedupScopeID(), event.getTsFileEpoch().getFilePath());
+ clearTsFileEpoch(event.getTsFileEpoch());
+ return null;
+ }
+
+ if (event.increaseReferenceCount(PipeRealtimeDataRegionHybridSource.class.getName())) {
+ if (isRegionLevelDowngradingEnabled) {
+ ++inFlightTsFileCount;
+ final PipeTsFileInsertionEvent tsFileInsertionEvent =
+ (PipeTsFileInsertionEvent) event.getEvent();
+ final Runnable clearTsFileEpochHook =
+ () -> clearTsFileEpochAfterCommit(event.getTsFileEpoch());
+ tsFileInsertionEvent.addOnTransferredHook(clearTsFileEpochHook);
+ tsFileInsertionEvent.addOnDiscardedHook(clearTsFileEpochHook);
} else {
- // If the event's reference count can not be increased, it means the data represented by
- // this event is not reliable anymore. the data has been lost. we simply discard this
- // event and report the exception to PipeRuntimeAgent.
- final String errorMessage =
- String.format(
- DataNodePipeMessages.EVENT_CAN_NOT_BE_SUPPLIED_BECAUSE_DATA_IS_LOST,
- event.getEvent());
- LOGGER.error(errorMessage);
- PipeDataNodeAgent.runtime()
- .report(pipeTaskMeta, new PipeRuntimeNonCriticalException(errorMessage));
- PipeTsFileEpochProgressIndexKeeper.getInstance()
- .eliminateProgressIndex(
- dataRegionId, getTsFileDedupScopeID(), event.getTsFileEpoch().getFilePath());
- return null;
+ clearTsFileEpoch(event.getTsFileEpoch());
}
- } finally {
+ return event.getEvent();
+ } else {
+ // If the event's reference count can not be increased, it means the data represented by
+ // this event is not reliable anymore. the data has been lost. we simply discard this
+ // event and report the exception to PipeRuntimeAgent.
+ final String errorMessage =
+ String.format(
+ DataNodePipeMessages.EVENT_CAN_NOT_BE_SUPPLIED_BECAUSE_DATA_IS_LOST,
+ event.getEvent());
+ LOGGER.error(errorMessage);
+ PipeDataNodeAgent.runtime()
+ .report(pipeTaskMeta, new PipeRuntimeNonCriticalException(errorMessage));
+ PipeTsFileEpochProgressIndexKeeper.getInstance()
+ .eliminateProgressIndex(
+ dataRegionId, getTsFileDedupScopeID(), event.getTsFileEpoch().getFilePath());
clearTsFileEpoch(event.getTsFileEpoch());
+ return null;
}
}
}
diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/PipeRealtimeExtractTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/PipeRealtimeExtractTest.java
index fc4a15a885b5e..6704aa245c5af 100644
--- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/PipeRealtimeExtractTest.java
+++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/PipeRealtimeExtractTest.java
@@ -19,6 +19,7 @@
package org.apache.iotdb.db.pipe.source;
+import org.apache.iotdb.commons.conf.CommonDescriptor;
import org.apache.iotdb.commons.conf.IoTDBConstant;
import org.apache.iotdb.commons.consensus.index.impl.MinimumProgressIndex;
import org.apache.iotdb.commons.path.PartialPath;
@@ -29,6 +30,7 @@
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStaticMeta;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta;
import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTemporaryMetaInAgent;
+import org.apache.iotdb.commons.pipe.agent.task.progress.PipeEventCommitManager;
import org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant;
import org.apache.iotdb.commons.pipe.config.plugin.configuraion.PipeTaskRuntimeConfiguration;
import org.apache.iotdb.commons.pipe.config.plugin.env.PipeTaskSourceRuntimeEnvironment;
@@ -38,12 +40,16 @@
import org.apache.iotdb.commons.utils.FileUtils;
import org.apache.iotdb.db.conf.IoTDBDescriptor;
import org.apache.iotdb.db.pipe.agent.PipeDataNodeAgent;
+import org.apache.iotdb.db.pipe.event.common.tablet.PipeRawTabletInsertionEvent;
+import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent;
import org.apache.iotdb.db.pipe.event.realtime.PipeRealtimeEvent;
import org.apache.iotdb.db.pipe.event.realtime.PipeRealtimeEventFactory;
import org.apache.iotdb.db.pipe.source.dataregion.realtime.PipeRealtimeDataRegionHybridSource;
import org.apache.iotdb.db.pipe.source.dataregion.realtime.PipeRealtimeDataRegionLogSource;
import org.apache.iotdb.db.pipe.source.dataregion.realtime.PipeRealtimeDataRegionSource;
import org.apache.iotdb.db.pipe.source.dataregion.realtime.PipeRealtimeDataRegionTsFileSource;
+import org.apache.iotdb.db.pipe.source.dataregion.realtime.assigner.PipeTsFileEpochProgressIndexKeeper;
+import org.apache.iotdb.db.pipe.source.dataregion.realtime.epoch.TsFileEpoch;
import org.apache.iotdb.db.pipe.source.dataregion.realtime.listener.PipeInsertionDataNodeListener;
import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertRowNode;
import org.apache.iotdb.db.storageengine.dataregion.tsfile.TsFileResource;
@@ -56,6 +62,8 @@
import org.apache.tsfile.common.constant.TsFileConstant;
import org.apache.tsfile.enums.TSDataType;
import org.apache.tsfile.file.metadata.IDeviceID;
+import org.apache.tsfile.write.record.Tablet;
+import org.apache.tsfile.write.schema.MeasurementSchema;
import org.junit.After;
import org.junit.Assert;
import org.junit.Before;
@@ -68,8 +76,10 @@
import java.lang.reflect.Field;
import java.nio.file.Files;
import java.util.Arrays;
+import java.util.Collections;
import java.util.HashMap;
import java.util.List;
+import java.util.Set;
import java.util.concurrent.ExecutionException;
import java.util.concurrent.ExecutorService;
import java.util.concurrent.Executors;
@@ -99,11 +109,14 @@ public class PipeRealtimeExtractTest {
private ExecutorService writeService;
private ExecutorService listenerService;
private int dataNodeId;
+ private double pipeTotalFloatingMemoryProportion;
@Before
public void setUp() throws Exception {
dataNodeId = IoTDBDescriptor.getInstance().getConfig().getDataNodeId();
IoTDBDescriptor.getInstance().getConfig().setDataNodeId(0);
+ pipeTotalFloatingMemoryProportion =
+ CommonDescriptor.getInstance().getConfig().getPipeTotalFloatingMemoryProportion();
removeTestPipeMeta();
writeService = Executors.newFixedThreadPool(2);
listenerService = Executors.newFixedThreadPool(4);
@@ -120,6 +133,9 @@ public void setUp() throws Exception {
@After
public void tearDown() throws Exception {
IoTDBDescriptor.getInstance().getConfig().setDataNodeId(dataNodeId);
+ CommonDescriptor.getInstance()
+ .getConfig()
+ .setPipeTotalFloatingMemoryProportion(pipeTotalFloatingMemoryProportion);
writeService.shutdownNow();
listenerService.shutdownNow();
FileUtils.deleteFileOrDirectory(tmpDir);
@@ -379,6 +395,782 @@ public void testHybridSourceReportsTsFileEpochDegradedStatus() throws Exception
Assert.assertNull(getGlobalTsFileEpochDegraded());
}
+ @Test
+ public void testHybridSourceRegionLevelDowngradingIsPipeSpecific() throws Exception {
+ try (final PipeRealtimeDataRegionHybridSource disabledExtractor =
+ new PipeRealtimeDataRegionHybridSource();
+ final PipeRealtimeDataRegionHybridSource enabledExtractor =
+ new PipeRealtimeDataRegionHybridSource()) {
+ final PipeParameters disabledParameters =
+ new PipeParameters(
+ new HashMap() {
+ {
+ put(PipeSourceConstant.EXTRACTOR_PATTERN_KEY, pattern1);
+ }
+ });
+ final PipeParameters enabledParameters =
+ new PipeParameters(
+ new HashMap() {
+ {
+ put(PipeSourceConstant.EXTRACTOR_PATTERN_KEY, pattern1);
+ put(
+ PipeSourceConstant.EXTRACTOR_REALTIME_REGION_LEVEL_DOWNGRADING_KEY,
+ Boolean.TRUE.toString());
+ }
+ });
+
+ final PipeTaskRuntimeConfiguration disabledConfiguration =
+ new PipeTaskRuntimeConfiguration(
+ new PipeTaskSourceRuntimeEnvironment(
+ "region-level-downgrading-disabled",
+ TEST_PIPE_CREATION_TIME,
+ dataRegion1,
+ new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 1)));
+ final PipeTaskRuntimeConfiguration enabledConfiguration =
+ new PipeTaskRuntimeConfiguration(
+ new PipeTaskSourceRuntimeEnvironment(
+ "region-level-downgrading-enabled",
+ TEST_PIPE_CREATION_TIME,
+ dataRegion1,
+ new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 1)));
+
+ disabledExtractor.validate(new PipeParameterValidator(disabledParameters));
+ disabledExtractor.customize(disabledParameters, disabledConfiguration);
+ enabledExtractor.validate(new PipeParameterValidator(enabledParameters));
+ enabledExtractor.customize(enabledParameters, enabledConfiguration);
+
+ Assert.assertFalse(isRegionLevelDowngradingEnabled(disabledExtractor));
+ Assert.assertTrue(isRegionLevelDowngradingEnabled(enabledExtractor));
+ }
+ }
+
+ @Test
+ public void testHybridSourceRegionLevelDowngradingWaitsForTsFileCommit() throws Exception {
+ registerTestPipeMeta();
+
+ final PipeEventCommitManager commitManager = PipeEventCommitManager.getInstance();
+ commitManager.register(TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1, "test");
+ try (final PipeRealtimeDataRegionHybridSource extractor =
+ new PipeRealtimeDataRegionHybridSource()) {
+ final PipeParameters parameters =
+ new PipeParameters(
+ new HashMap() {
+ {
+ put(PipeSourceConstant.EXTRACTOR_PATTERN_KEY, pattern1);
+ put(
+ PipeSourceConstant.SOURCE_REALTIME_REGION_LEVEL_DOWNGRADING_KEY,
+ Boolean.TRUE.toString());
+ }
+ });
+ final PipeTaskMeta pipeTaskMeta = new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 1);
+ final PipeTaskRuntimeConfiguration configuration =
+ new PipeTaskRuntimeConfiguration(
+ new PipeTaskSourceRuntimeEnvironment(
+ TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1, pipeTaskMeta));
+
+ extractor.validate(new PipeParameterValidator(parameters));
+ extractor.customize(parameters, configuration);
+
+ final TsFileResource firstResource = createTsFileResource(dataRegion1, "101-101-0-0.tsfile");
+ final PipeRealtimeEvent firstTabletEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(
+ false,
+ "root.sg",
+ createInsertRowNode("first-degraded-tablet", "a"),
+ firstResource),
+ extractor,
+ pipeTaskMeta);
+
+ Assert.assertTrue(firstTabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(firstTabletEvent);
+ Assert.assertEquals(Boolean.FALSE, getGlobalTsFileEpochDegraded());
+
+ firstTabletEvent.clearReferenceCount(TEST_REFERENCE_HOLDER);
+
+ // Queue a tablet from another epoch before the first epoch triggers region-level
+ // downgrading. It should be buffered while the degraded TsFile is being sent.
+ final TsFileResource secondResource = createTsFileResource(dataRegion1, "102-102-0-0.tsfile");
+ final PipeRealtimeEvent secondTabletEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(
+ false,
+ "root.sg",
+ createInsertRowNode("second-degraded-tablet", "a"),
+ secondResource),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(secondTabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(secondTabletEvent);
+ Assert.assertEquals(
+ TsFileEpoch.State.USING_TABLET, secondTabletEvent.getTsFileEpoch().getState(extractor));
+
+ Assert.assertNull(extractor.supply());
+ Assert.assertEquals(Boolean.TRUE, getGlobalTsFileEpochDegraded());
+ Assert.assertEquals(
+ TsFileEpoch.State.USING_TABLET, secondTabletEvent.getTsFileEpoch().getState(extractor));
+ Assert.assertFalse(secondTabletEvent.getEvent().isReleased());
+
+ // Simulate that the buffered tablet is evicted before the previous degraded TsFile is
+ // committed. The second epoch should then continue region-level downgrading with its TsFile.
+ secondTabletEvent.clearReferenceCount(TEST_REFERENCE_HOLDER);
+
+ final PipeRealtimeEvent firstTsFileEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(false, "root.sg", firstResource, false),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(firstTsFileEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(firstTsFileEvent);
+ Assert.assertEquals(
+ TsFileEpoch.State.USING_TSFILE, firstTsFileEvent.getTsFileEpoch().getState(extractor));
+
+ final Event firstSuppliedTsFile = extractor.supply();
+ Assert.assertTrue(firstSuppliedTsFile instanceof TsFileInsertionEvent);
+ Assert.assertEquals(Boolean.TRUE, getGlobalTsFileEpochDegraded());
+
+ final PipeRealtimeEvent secondTsFileEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(false, "root.sg", secondResource, false),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(secondTsFileEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(secondTsFileEvent);
+
+ // The second TsFile stays in the source until the first TsFile is committed downstream.
+ Assert.assertNull(extractor.supply());
+ final PipeTsFileInsertionEvent suppliedFirstTsFile =
+ (PipeTsFileInsertionEvent) firstSuppliedTsFile;
+ suppliedFirstTsFile.registerGeneratedTabletInsertionEvent();
+ suppliedFirstTsFile.registerGeneratedTabletInsertionEvent();
+ suppliedFirstTsFile.markGeneratedTabletInsertionEventsParsingCompleted();
+ final PipeRawTabletInsertionEvent firstGeneratedTabletEvent =
+ createGeneratedTabletEvent(suppliedFirstTsFile, pipeTaskMeta, "first-generated");
+ final PipeRawTabletInsertionEvent secondGeneratedTabletEvent =
+ createGeneratedTabletEvent(suppliedFirstTsFile, pipeTaskMeta, "second-generated");
+ Assert.assertTrue(firstGeneratedTabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ Assert.assertTrue(secondGeneratedTabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ Assert.assertTrue(suppliedFirstTsFile.decreaseReferenceCount(TEST_REFERENCE_HOLDER, false));
+ commitSuppliedEvent(firstGeneratedTabletEvent, commitManager);
+ Assert.assertEquals(Boolean.TRUE, getGlobalTsFileEpochDegraded());
+ Assert.assertEquals(2, getActiveTsFileEpochCount(extractor));
+ Assert.assertEquals(1, getInFlightTsFileCount(extractor));
+ commitSuppliedEvent(secondGeneratedTabletEvent, commitManager);
+ Assert.assertEquals(1, getActiveTsFileEpochCount(extractor));
+ Assert.assertEquals(0, getInFlightTsFileCount(extractor));
+ Assert.assertEquals(Boolean.TRUE, getGlobalTsFileEpochDegraded());
+
+ final Event secondSuppliedTsFile = extractor.supply();
+ Assert.assertTrue(secondSuppliedTsFile instanceof TsFileInsertionEvent);
+ Assert.assertEquals(Boolean.TRUE, getGlobalTsFileEpochDegraded());
+
+ commitSuppliedEvent(secondSuppliedTsFile, commitManager);
+ Assert.assertEquals(0, getActiveTsFileEpochCount(extractor));
+ Assert.assertEquals(0, getInFlightTsFileCount(extractor));
+ Assert.assertNull(getGlobalTsFileEpochDegraded());
+ Assert.assertNull(extractor.supply());
+ } finally {
+ commitManager.deregister(TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1);
+ }
+ }
+
+ @Test
+ public void testGeneratedTabletTransferWaitsForAllTabletCommits() throws Exception {
+ registerTestPipeMeta();
+
+ final PipeEventCommitManager commitManager = PipeEventCommitManager.getInstance();
+ final String dedupScopeId = "generated-tablet-transfer-test";
+ commitManager.register(TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1, "test");
+ try {
+ final PipeTaskMeta pipeTaskMeta = new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 1);
+ final TsFileResource resource = createTsFileResource(dataRegion1, "110-110-0-0.tsfile");
+ final PipeTsFileInsertionEvent tsFileEvent =
+ new PipeTsFileInsertionEvent(false, "root.sg", resource, false)
+ .shallowCopySelfAndBindPipeTaskMetaForProgressReport(
+ TEST_PIPE_NAME,
+ TEST_PIPE_CREATION_TIME,
+ pipeTaskMeta,
+ null,
+ null,
+ null,
+ null,
+ null,
+ true,
+ Long.MIN_VALUE,
+ Long.MAX_VALUE);
+ tsFileEvent.bindTsFileDedupScopeID(dedupScopeId);
+ PipeTsFileEpochProgressIndexKeeper.getInstance()
+ .registerProgressIndex(dataRegion1, dedupScopeId, resource);
+
+ final AtomicBoolean transferred = new AtomicBoolean(false);
+ tsFileEvent.addOnTransferredHook(() -> transferred.set(true));
+ tsFileEvent.registerGeneratedTabletInsertionEvent();
+ tsFileEvent.registerGeneratedTabletInsertionEvent();
+
+ // The default PipeProcessor path iterates toTabletInsertionEvents() directly and does not
+ // report parser completion. The source TsFile commit is still the boundary before generated
+ // tablet commits.
+ tsFileEvent.skipReportOnCommit();
+ tsFileEvent.getOnCommittedHooks().forEach(Runnable::run);
+
+ final PipeRawTabletInsertionEvent firstGeneratedTabletEvent =
+ createGeneratedTabletEvent(tsFileEvent, pipeTaskMeta, "first", false);
+ final PipeRawTabletInsertionEvent secondGeneratedTabletEvent =
+ createGeneratedTabletEvent(tsFileEvent, pipeTaskMeta, "second", true);
+ Assert.assertTrue(firstGeneratedTabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ Assert.assertTrue(secondGeneratedTabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+
+ commitSuppliedEvent(firstGeneratedTabletEvent, commitManager);
+ Assert.assertFalse(transferred.get());
+ Assert.assertTrue(
+ PipeTsFileEpochProgressIndexKeeper.getInstance()
+ .containsTsFile(dataRegion1, dedupScopeId, resource.getTsFilePath()));
+
+ commitSuppliedEvent(secondGeneratedTabletEvent, commitManager);
+ Assert.assertTrue(transferred.get());
+ Assert.assertFalse(
+ PipeTsFileEpochProgressIndexKeeper.getInstance()
+ .containsTsFile(dataRegion1, dedupScopeId, resource.getTsFilePath()));
+ } finally {
+ commitManager.deregister(TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1);
+ PipeTsFileEpochProgressIndexKeeper.getInstance()
+ .clearProgressIndex(dataRegion1, dedupScopeId);
+ }
+ }
+
+ @Test
+ public void testGeneratedTabletTransferWaitsForDeferredGeneration() throws Exception {
+ registerTestPipeMeta();
+
+ final PipeEventCommitManager commitManager = PipeEventCommitManager.getInstance();
+ commitManager.register(TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1, "test");
+ final String dedupScopeId = "deferred-generated-tablet-transfer-test";
+ try {
+ final PipeTaskMeta pipeTaskMeta = new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 1);
+ final TsFileResource resource = createTsFileResource(dataRegion1, "111-111-0-0.tsfile");
+ final PipeTsFileInsertionEvent tsFileEvent =
+ new PipeTsFileInsertionEvent(false, "root.sg", resource, false)
+ .shallowCopySelfAndBindPipeTaskMetaForProgressReport(
+ TEST_PIPE_NAME,
+ TEST_PIPE_CREATION_TIME,
+ pipeTaskMeta,
+ null,
+ null,
+ null,
+ null,
+ null,
+ true,
+ Long.MIN_VALUE,
+ Long.MAX_VALUE);
+ tsFileEvent.bindTsFileDedupScopeID(dedupScopeId);
+
+ final AtomicBoolean transferred = new AtomicBoolean(false);
+ tsFileEvent.addOnTransferredHook(() -> transferred.set(true));
+ tsFileEvent.markGeneratedTabletInsertionEventsParsingStarted();
+ tsFileEvent.skipReportOnCommit();
+ tsFileEvent.getOnCommittedHooks().forEach(Runnable::run);
+ Assert.assertFalse(transferred.get());
+
+ tsFileEvent.registerGeneratedTabletInsertionEvent();
+ tsFileEvent.markGeneratedTabletInsertionEventsParsingCompleted();
+ final PipeRawTabletInsertionEvent generatedTabletEvent =
+ createGeneratedTabletEvent(tsFileEvent, pipeTaskMeta, "deferred");
+ Assert.assertTrue(generatedTabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ commitSuppliedEvent(generatedTabletEvent, commitManager);
+ Assert.assertTrue(transferred.get());
+ } finally {
+ commitManager.deregister(TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1);
+ PipeTsFileEpochProgressIndexKeeper.getInstance()
+ .clearProgressIndex(dataRegion1, dedupScopeId);
+ }
+ }
+
+ @Test
+ public void testHybridSourceClearsInFlightTsFileWhenSuppliedEventIsDiscarded() throws Exception {
+ registerTestPipeMeta();
+
+ final PipeEventCommitManager commitManager = PipeEventCommitManager.getInstance();
+ commitManager.register(TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1, "test");
+ try (final PipeRealtimeDataRegionHybridSource extractor =
+ new PipeRealtimeDataRegionHybridSource()) {
+ final PipeTaskMeta pipeTaskMeta = new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 1);
+ final PipeTaskRuntimeConfiguration configuration =
+ new PipeTaskRuntimeConfiguration(
+ new PipeTaskSourceRuntimeEnvironment(
+ TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1, pipeTaskMeta));
+ final PipeParameters parameters =
+ new PipeParameters(
+ new HashMap() {
+ {
+ put(PipeSourceConstant.EXTRACTOR_PATTERN_KEY, pattern1);
+ put(
+ PipeSourceConstant.SOURCE_REALTIME_REGION_LEVEL_DOWNGRADING_KEY,
+ Boolean.TRUE.toString());
+ }
+ });
+ extractor.validate(new PipeParameterValidator(parameters));
+ extractor.customize(parameters, configuration);
+
+ final TsFileResource resource = createTsFileResource(dataRegion1, "112-112-0-0.tsfile");
+ final PipeRealtimeEvent tabletEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(
+ false, "root.sg", createInsertRowNode("discarded-tsfile-tablet", "a"), resource),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(tabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(tabletEvent);
+ tabletEvent.clearReferenceCount(TEST_REFERENCE_HOLDER);
+ Assert.assertNull(extractor.supply());
+ Assert.assertEquals(Boolean.TRUE, getGlobalTsFileEpochDegraded());
+
+ final PipeRealtimeEvent tsFileEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(false, "root.sg", resource, false),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(tsFileEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(tsFileEvent);
+
+ final Event suppliedTsFile = extractor.supply();
+ Assert.assertTrue(suppliedTsFile instanceof TsFileInsertionEvent);
+ Assert.assertEquals(1, getInFlightTsFileCount(extractor));
+
+ ((EnrichedEvent) suppliedTsFile).clearReferenceCount(TEST_REFERENCE_HOLDER);
+ Assert.assertEquals(0, getInFlightTsFileCount(extractor));
+ Assert.assertNull(getGlobalTsFileEpochDegraded());
+ } finally {
+ commitManager.deregister(TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1);
+ }
+ }
+
+ @Test
+ public void testHybridSourceCompensatesForDiscardedGeneratedTabletEvents() throws Exception {
+ registerTestPipeMeta();
+
+ final PipeEventCommitManager commitManager = PipeEventCommitManager.getInstance();
+ commitManager.register(TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1, "test");
+ try (final PipeRealtimeDataRegionHybridSource extractor =
+ new PipeRealtimeDataRegionHybridSource()) {
+ final PipeTaskMeta pipeTaskMeta = new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 1);
+ final PipeTaskRuntimeConfiguration configuration =
+ new PipeTaskRuntimeConfiguration(
+ new PipeTaskSourceRuntimeEnvironment(
+ TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1, pipeTaskMeta));
+ final PipeParameters parameters =
+ new PipeParameters(
+ new HashMap() {
+ {
+ put(PipeSourceConstant.EXTRACTOR_PATTERN_KEY, pattern1);
+ put(
+ PipeSourceConstant.SOURCE_REALTIME_REGION_LEVEL_DOWNGRADING_KEY,
+ Boolean.TRUE.toString());
+ }
+ });
+ extractor.validate(new PipeParameterValidator(parameters));
+ extractor.customize(parameters, configuration);
+
+ final TsFileResource resource = createTsFileResource(dataRegion1, "113-113-0-0.tsfile");
+ final PipeRealtimeEvent tabletEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(
+ false,
+ "root.sg",
+ createInsertRowNode("discarded-generated-tablet", "a"),
+ resource),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(tabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(tabletEvent);
+ tabletEvent.clearReferenceCount(TEST_REFERENCE_HOLDER);
+ Assert.assertNull(extractor.supply());
+
+ final PipeRealtimeEvent tsFileRealtimeEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(false, "root.sg", resource, false),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(tsFileRealtimeEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(tsFileRealtimeEvent);
+ final PipeTsFileInsertionEvent suppliedTsFile = (PipeTsFileInsertionEvent) extractor.supply();
+ Assert.assertNotNull(suppliedTsFile);
+ Assert.assertEquals(1, getInFlightTsFileCount(extractor));
+
+ suppliedTsFile.registerGeneratedTabletInsertionEvent();
+ suppliedTsFile.registerGeneratedTabletInsertionEvent();
+ suppliedTsFile.markGeneratedTabletInsertionEventsParsingCompleted();
+ final PipeRawTabletInsertionEvent firstGeneratedTablet =
+ createGeneratedTabletEvent(suppliedTsFile, pipeTaskMeta, "discarded-first");
+ final PipeRawTabletInsertionEvent secondGeneratedTablet =
+ createGeneratedTabletEvent(suppliedTsFile, pipeTaskMeta, "discarded-second");
+ firstGeneratedTablet.markAsGeneratedEventRegisteredWithSource();
+ secondGeneratedTablet.markAsGeneratedEventRegisteredWithSource();
+ Assert.assertTrue(firstGeneratedTablet.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ Assert.assertTrue(secondGeneratedTablet.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+
+ commitManager.enrichWithCommitterKeyAndCommitId(
+ suppliedTsFile, TEST_PIPE_CREATION_TIME, dataRegion1);
+ Assert.assertTrue(suppliedTsFile.decreaseReferenceCount(TEST_REFERENCE_HOLDER, false));
+ firstGeneratedTablet.clearReferenceCount(TEST_REFERENCE_HOLDER);
+ Assert.assertEquals(1, getInFlightTsFileCount(extractor));
+ secondGeneratedTablet.clearReferenceCount(TEST_REFERENCE_HOLDER);
+
+ Assert.assertEquals(0, getInFlightTsFileCount(extractor));
+ Assert.assertNull(getGlobalTsFileEpochDegraded());
+ } finally {
+ commitManager.deregister(TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1);
+ }
+ }
+
+ @Test
+ public void testHybridSourceRegionLevelDowngradingResumesCompleteBufferedTablets()
+ throws Exception {
+ registerTestPipeMeta();
+
+ final PipeEventCommitManager commitManager = PipeEventCommitManager.getInstance();
+ commitManager.register(TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1, "test");
+ try (final PipeRealtimeDataRegionHybridSource extractor =
+ new PipeRealtimeDataRegionHybridSource()) {
+ final PipeParameters parameters =
+ new PipeParameters(
+ new HashMap() {
+ {
+ put(PipeSourceConstant.EXTRACTOR_PATTERN_KEY, pattern1);
+ put(
+ PipeSourceConstant.SOURCE_REALTIME_REGION_LEVEL_DOWNGRADING_KEY,
+ Boolean.TRUE.toString());
+ }
+ });
+ final PipeTaskMeta pipeTaskMeta = new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 1);
+ final PipeTaskRuntimeConfiguration configuration =
+ new PipeTaskRuntimeConfiguration(
+ new PipeTaskSourceRuntimeEnvironment(
+ TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1, pipeTaskMeta));
+
+ extractor.validate(new PipeParameterValidator(parameters));
+ extractor.customize(parameters, configuration);
+
+ final TsFileResource firstResource = createTsFileResource(dataRegion1, "103-103-0-0.tsfile");
+ final PipeRealtimeEvent firstTabletEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(
+ false,
+ "root.sg",
+ createInsertRowNode("first-degraded-tablet", "a"),
+ firstResource),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(firstTabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(firstTabletEvent);
+ firstTabletEvent.clearReferenceCount(TEST_REFERENCE_HOLDER);
+
+ final TsFileResource secondResource = createTsFileResource(dataRegion1, "104-104-0-0.tsfile");
+ final PipeRealtimeEvent secondTabletEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(
+ false,
+ "root.sg",
+ createInsertRowNode("fully-buffered-tablet", "a"),
+ secondResource),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(secondTabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(secondTabletEvent);
+
+ Assert.assertNull(extractor.supply());
+ Assert.assertEquals(Boolean.TRUE, getGlobalTsFileEpochDegraded());
+ Assert.assertEquals(
+ TsFileEpoch.State.USING_TABLET, secondTabletEvent.getTsFileEpoch().getState(extractor));
+ Assert.assertFalse(secondTabletEvent.getEvent().isReleased());
+
+ final PipeRealtimeEvent firstTsFileEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(false, "root.sg", firstResource, false),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(firstTsFileEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(firstTsFileEvent);
+ final Event firstSuppliedTsFile = extractor.supply();
+ Assert.assertTrue(firstSuppliedTsFile instanceof TsFileInsertionEvent);
+
+ commitLastGeneratedTabletEvent(
+ (PipeTsFileInsertionEvent) firstSuppliedTsFile, commitManager, pipeTaskMeta);
+ Assert.assertEquals(Boolean.FALSE, getGlobalTsFileEpochDegraded());
+
+ // The latest TsFile is still open. Since all of its requests survived in memory at the
+ // commit boundary above, later writes of the same TsFile should immediately continue on the
+ // realtime path instead of waiting for another flush.
+ final PipeRealtimeEvent newRealtimeTabletEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(
+ false,
+ "root.sg",
+ createInsertRowNode("new-realtime-tablet", "a"),
+ secondResource),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(newRealtimeTabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(newRealtimeTabletEvent);
+
+ final Event resumedTabletEvent = extractor.supply();
+ Assert.assertTrue(resumedTabletEvent instanceof TabletInsertionEvent);
+ Assert.assertSame(secondTabletEvent.getEvent(), resumedTabletEvent);
+ Assert.assertEquals(Boolean.FALSE, getGlobalTsFileEpochDegraded());
+ commitSuppliedEvent(resumedTabletEvent, commitManager);
+
+ final Event newSuppliedTabletEvent = extractor.supply();
+ Assert.assertTrue(newSuppliedTabletEvent instanceof TabletInsertionEvent);
+ Assert.assertSame(newRealtimeTabletEvent.getEvent(), newSuppliedTabletEvent);
+ Assert.assertEquals(Boolean.FALSE, getGlobalTsFileEpochDegraded());
+ commitSuppliedEvent(newSuppliedTabletEvent, commitManager);
+
+ final PipeRealtimeEvent secondTsFileEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(false, "root.sg", secondResource, false),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(secondTsFileEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(secondTsFileEvent);
+
+ // The second TsFile is no longer needed because all of its tablets survived buffering.
+ Assert.assertNull(extractor.supply());
+ Assert.assertNull(getGlobalTsFileEpochDegraded());
+ } finally {
+ commitManager.deregister(TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1);
+ }
+ }
+
+ @Test
+ public void testHybridSourceRegionLevelDowngradingOnlyCachesLatestTsFile() throws Exception {
+ registerTestPipeMeta();
+
+ final PipeEventCommitManager commitManager = PipeEventCommitManager.getInstance();
+ commitManager.register(TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1, "test");
+ try (final PipeRealtimeDataRegionHybridSource extractor =
+ new PipeRealtimeDataRegionHybridSource()) {
+ final PipeParameters parameters =
+ new PipeParameters(
+ new HashMap() {
+ {
+ put(PipeSourceConstant.EXTRACTOR_PATTERN_KEY, pattern1);
+ put(
+ PipeSourceConstant.SOURCE_REALTIME_REGION_LEVEL_DOWNGRADING_KEY,
+ Boolean.TRUE.toString());
+ }
+ });
+ final PipeTaskMeta pipeTaskMeta = new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 1);
+ final PipeTaskRuntimeConfiguration configuration =
+ new PipeTaskRuntimeConfiguration(
+ new PipeTaskSourceRuntimeEnvironment(
+ TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1, pipeTaskMeta));
+
+ extractor.validate(new PipeParameterValidator(parameters));
+ extractor.customize(parameters, configuration);
+
+ final TsFileResource firstResource = createTsFileResource(dataRegion1, "107-107-0-0.tsfile");
+ final PipeRealtimeEvent firstTabletEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(
+ false,
+ "root.sg",
+ createInsertRowNode("first-degraded-tablet", "a"),
+ firstResource),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(firstTabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(firstTabletEvent);
+ firstTabletEvent.clearReferenceCount(TEST_REFERENCE_HOLDER);
+
+ Assert.assertNull(extractor.supply());
+ Assert.assertEquals(Boolean.TRUE, getGlobalTsFileEpochDegraded());
+
+ final TsFileResource secondResource = createTsFileResource(dataRegion1, "108-108-0-0.tsfile");
+ final PipeRealtimeEvent secondTabletEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(
+ false,
+ "root.sg",
+ createInsertRowNode("second-buffered-tablet", "a"),
+ secondResource),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(secondTabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(secondTabletEvent);
+
+ final PipeRealtimeEvent secondTsFileEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(false, "root.sg", secondResource, false),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(secondTsFileEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(secondTsFileEvent);
+
+ final TsFileResource thirdResource = createTsFileResource(dataRegion1, "109-109-0-0.tsfile");
+ final PipeRealtimeEvent thirdTabletEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(
+ false,
+ "root.sg",
+ createInsertRowNode("latest-buffered-tablet", "a"),
+ thirdResource),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(thirdTabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(thirdTabletEvent);
+
+ // Once a newer epoch appears, the former tail is downgraded even if all its tablets are
+ // still available. This bounds the region-level cache to the latest TsFile.
+ Assert.assertEquals(
+ TsFileEpoch.State.USING_TSFILE, secondTabletEvent.getTsFileEpoch().getState(extractor));
+ Assert.assertTrue(secondTabletEvent.getEvent().isReleased());
+ Assert.assertEquals(
+ TsFileEpoch.State.USING_TABLET, thirdTabletEvent.getTsFileEpoch().getState(extractor));
+ Assert.assertFalse(thirdTabletEvent.getEvent().isReleased());
+
+ // Extract the first TsFile after the second one to verify that downgrade order, rather than
+ // flush completion order, decides which file can pass downstream.
+ final PipeRealtimeEvent firstTsFileEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(false, "root.sg", firstResource, false),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(firstTsFileEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(firstTsFileEvent);
+
+ final Event firstSuppliedTsFile = extractor.supply();
+ Assert.assertTrue(firstSuppliedTsFile instanceof TsFileInsertionEvent);
+ Assert.assertSame(firstTsFileEvent.getEvent(), firstSuppliedTsFile);
+ commitSuppliedEvent(firstSuppliedTsFile, commitManager);
+
+ final Event secondSuppliedTsFile = extractor.supply();
+ Assert.assertTrue(secondSuppliedTsFile instanceof TsFileInsertionEvent);
+ Assert.assertSame(secondTsFileEvent.getEvent(), secondSuppliedTsFile);
+ commitSuppliedEvent(secondSuppliedTsFile, commitManager);
+
+ Assert.assertEquals(Boolean.FALSE, getGlobalTsFileEpochDegraded());
+ final Event resumedLatestTablet = extractor.supply();
+ Assert.assertTrue(resumedLatestTablet instanceof TabletInsertionEvent);
+ Assert.assertSame(thirdTabletEvent.getEvent(), resumedLatestTablet);
+ commitSuppliedEvent(resumedLatestTablet, commitManager);
+
+ final PipeRealtimeEvent thirdTsFileEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(false, "root.sg", thirdResource, false),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(thirdTsFileEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(thirdTsFileEvent);
+
+ Assert.assertNull(extractor.supply());
+ Assert.assertNull(getGlobalTsFileEpochDegraded());
+ } finally {
+ commitManager.deregister(TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1);
+ }
+ }
+
+ @Test
+ public void testHybridSourceRegionLevelDowngradingPreservesPreviouslyQueuedEvents()
+ throws Exception {
+ registerTestPipeMeta();
+
+ final PipeEventCommitManager commitManager = PipeEventCommitManager.getInstance();
+ commitManager.register(TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1, "test");
+ try (final PipeRealtimeDataRegionHybridSource extractor =
+ new PipeRealtimeDataRegionHybridSource()) {
+ final PipeParameters parameters =
+ new PipeParameters(
+ new HashMap() {
+ {
+ put(PipeSourceConstant.EXTRACTOR_PATTERN_KEY, pattern1);
+ put(
+ PipeSourceConstant.SOURCE_REALTIME_REGION_LEVEL_DOWNGRADING_KEY,
+ Boolean.TRUE.toString());
+ }
+ });
+ final PipeTaskMeta pipeTaskMeta = new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 1);
+ final PipeTaskRuntimeConfiguration configuration =
+ new PipeTaskRuntimeConfiguration(
+ new PipeTaskSourceRuntimeEnvironment(
+ TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1, pipeTaskMeta));
+
+ extractor.validate(new PipeParameterValidator(parameters));
+ extractor.customize(parameters, configuration);
+
+ final TsFileResource olderResource = createTsFileResource(dataRegion1, "105-105-0-0.tsfile");
+ final PipeRealtimeEvent olderTabletEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(
+ false,
+ "root.sg",
+ createInsertRowNode("queued-before-downgrading", "a"),
+ olderResource),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(olderTabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(olderTabletEvent);
+
+ // Seal the older epoch while leaving its tablet queued in the source.
+ final PipeRealtimeEvent olderTsFileEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(false, "root.sg", olderResource, false),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(olderTsFileEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(olderTsFileEvent);
+
+ final TsFileResource degradedResource =
+ createTsFileResource(dataRegion1, "106-106-0-0.tsfile");
+ final PipeRealtimeEvent degradedTabletEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(
+ false,
+ "root.sg",
+ createInsertRowNode("trigger-region-downgrading", "a"),
+ degradedResource),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(degradedTabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+
+ CommonDescriptor.getInstance().getConfig().setPipeTotalFloatingMemoryProportion(0);
+ try {
+ extractor.extract(degradedTabletEvent);
+ } finally {
+ CommonDescriptor.getInstance()
+ .getConfig()
+ .setPipeTotalFloatingMemoryProportion(pipeTotalFloatingMemoryProportion);
+ }
+ Assert.assertEquals(
+ TsFileEpoch.State.USING_TSFILE, degradedTabletEvent.getTsFileEpoch().getState(extractor));
+ Assert.assertEquals(Boolean.TRUE, getGlobalTsFileEpochDegraded());
+
+ final PipeRealtimeEvent degradedTsFileEvent =
+ bindToTestPipe(
+ PipeRealtimeEventFactory.createRealtimeEvent(
+ false, "root.sg", degradedResource, false),
+ extractor,
+ pipeTaskMeta);
+ Assert.assertTrue(degradedTsFileEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ extractor.extract(degradedTsFileEvent);
+
+ // The tablet that was already queued before downgrading must not be overtaken by the later
+ // degraded TsFile.
+ final Event firstSuppliedEvent = extractor.supply();
+ Assert.assertTrue(firstSuppliedEvent instanceof TabletInsertionEvent);
+ Assert.assertSame(olderTabletEvent.getEvent(), firstSuppliedEvent);
+ commitSuppliedEvent(firstSuppliedEvent, commitManager);
+
+ final Event secondSuppliedEvent = extractor.supply();
+ Assert.assertTrue(secondSuppliedEvent instanceof TsFileInsertionEvent);
+ Assert.assertSame(degradedTsFileEvent.getEvent(), secondSuppliedEvent);
+ commitSuppliedEvent(secondSuppliedEvent, commitManager);
+
+ Assert.assertNull(getGlobalTsFileEpochDegraded());
+ Assert.assertNull(extractor.supply());
+ } finally {
+ commitManager.deregister(TEST_PIPE_NAME, TEST_PIPE_CREATION_TIME, dataRegion1);
+ }
+ }
+
private Future> write2DataRegion(
final int writeNum, final int dataRegionId, final int startNum) {
final File dataRegionDir =
@@ -577,6 +1369,101 @@ private void releaseSuppliedEvent(final Event event) {
}
}
+ private PipeRealtimeEvent bindToTestPipe(
+ final PipeRealtimeEvent event,
+ final PipeRealtimeDataRegionSource extractor,
+ final PipeTaskMeta pipeTaskMeta) {
+ return event.shallowCopySelfAndBindPipeTaskMetaForProgressReport(
+ TEST_PIPE_NAME,
+ TEST_PIPE_CREATION_TIME,
+ pipeTaskMeta,
+ extractor.getTreePattern(),
+ extractor.getTablePattern(),
+ String.valueOf(extractor.getUserId()),
+ extractor.getUserName(),
+ extractor.getCliHostname(),
+ extractor.isSkipIfNoPrivileges(),
+ extractor.getRealtimeDataExtractionStartTime(),
+ extractor.getRealtimeDataExtractionEndTime());
+ }
+
+ private void commitSuppliedEvent(final Event event, final PipeEventCommitManager commitManager) {
+ final EnrichedEvent enrichedEvent = (EnrichedEvent) event;
+ commitManager.enrichWithCommitterKeyAndCommitId(
+ enrichedEvent, TEST_PIPE_CREATION_TIME, dataRegion1);
+ Assert.assertTrue(enrichedEvent.decreaseReferenceCount(TEST_REFERENCE_HOLDER, true));
+ }
+
+ private void commitLastGeneratedTabletEvent(
+ final PipeTsFileInsertionEvent tsFileEvent,
+ final PipeEventCommitManager commitManager,
+ final PipeTaskMeta pipeTaskMeta) {
+ tsFileEvent.registerGeneratedTabletInsertionEvent();
+ tsFileEvent.markGeneratedTabletInsertionEventsParsingCompleted();
+ final PipeRawTabletInsertionEvent generatedTabletEvent =
+ createGeneratedTabletEvent(tsFileEvent, pipeTaskMeta, "generated");
+
+ Assert.assertTrue(generatedTabletEvent.increaseReferenceCount(TEST_REFERENCE_HOLDER));
+ Assert.assertTrue(tsFileEvent.decreaseReferenceCount(TEST_REFERENCE_HOLDER, false));
+ commitSuppliedEvent(generatedTabletEvent, commitManager);
+ }
+
+ private PipeRawTabletInsertionEvent createGeneratedTabletEvent(
+ final PipeTsFileInsertionEvent tsFileEvent,
+ final PipeTaskMeta pipeTaskMeta,
+ final String deviceId) {
+ return createGeneratedTabletEvent(tsFileEvent, pipeTaskMeta, deviceId, true);
+ }
+
+ private PipeRawTabletInsertionEvent createGeneratedTabletEvent(
+ final PipeTsFileInsertionEvent tsFileEvent,
+ final PipeTaskMeta pipeTaskMeta,
+ final String deviceId,
+ final boolean needToReport) {
+ final Tablet tablet =
+ new Tablet(
+ "root.sg.d." + deviceId,
+ Collections.singletonList(new MeasurementSchema("s", TSDataType.INT32)),
+ 1);
+ return new PipeRawTabletInsertionEvent(
+ false,
+ "root.sg",
+ null,
+ null,
+ tablet,
+ false,
+ TEST_PIPE_NAME,
+ TEST_PIPE_CREATION_TIME,
+ pipeTaskMeta,
+ tsFileEvent,
+ needToReport);
+ }
+
+ private int getActiveTsFileEpochCount(final PipeRealtimeDataRegionHybridSource extractor)
+ throws Exception {
+ final Field activeTsFileEpochsField =
+ PipeRealtimeDataRegionHybridSource.class.getDeclaredField("activeTsFileEpochs");
+ activeTsFileEpochsField.setAccessible(true);
+ return ((Set>) activeTsFileEpochsField.get(extractor)).size();
+ }
+
+ private int getInFlightTsFileCount(final PipeRealtimeDataRegionHybridSource extractor)
+ throws Exception {
+ final Field inFlightTsFileCountField =
+ PipeRealtimeDataRegionHybridSource.class.getDeclaredField("inFlightTsFileCount");
+ inFlightTsFileCountField.setAccessible(true);
+ return inFlightTsFileCountField.getInt(extractor);
+ }
+
+ private boolean isRegionLevelDowngradingEnabled(
+ final PipeRealtimeDataRegionHybridSource extractor) throws Exception {
+ final Field isRegionLevelDowngradingEnabledField =
+ PipeRealtimeDataRegionHybridSource.class.getDeclaredField(
+ "isRegionLevelDowngradingEnabled");
+ isRegionLevelDowngradingEnabledField.setAccessible(true);
+ return isRegionLevelDowngradingEnabledField.getBoolean(extractor);
+ }
+
private PipeRealtimeEvent createProgressReportRealtimeEvent() {
final ProgressReportEvent progressReportEvent = new ProgressReportEvent(null, 0, null);
progressReportEvent.bindProgressIndex(MinimumProgressIndex.INSTANCE);
diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSourceConstant.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSourceConstant.java
index 055553c7e1ab9..70d452ae36ace 100644
--- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSourceConstant.java
+++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/config/constant/PipeSourceConstant.java
@@ -142,6 +142,11 @@ public class PipeSourceConstant {
public static final String EXTRACTOR_REALTIME_LOOSE_RANGE_PATH_VALUE = "path";
public static final String EXTRACTOR_REALTIME_LOOSE_RANGE_ALL_VALUE = "all";
public static final String EXTRACTOR_REALTIME_LOOSE_RANGE_DEFAULT_VALUE = "";
+ public static final String EXTRACTOR_REALTIME_REGION_LEVEL_DOWNGRADING_KEY =
+ "extractor.realtime.region-level-downgrading";
+ public static final String SOURCE_REALTIME_REGION_LEVEL_DOWNGRADING_KEY =
+ "source.realtime.region-level-downgrading";
+ public static final boolean EXTRACTOR_REALTIME_REGION_LEVEL_DOWNGRADING_DEFAULT_VALUE = false;
public static final String EXTRACTOR_MODE_STREAMING_KEY = "extractor.mode.streaming";
public static final String SOURCE_MODE_STREAMING_KEY = "source.mode.streaming";