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 cc6b126384f6..d94b48092fe2 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 @@ -262,6 +262,13 @@ private void collectEvent(final Event event) { if (pendingQueue.offer(event)) { collectInvocationCount.incrementAndGet(); + if (event instanceof PipeRawTabletInsertionEvent + && ((PipeRawTabletInsertionEvent) event).getSourceEvent() + instanceof PipeTsFileInsertionEvent) { + ((PipeTsFileInsertionEvent) ((PipeRawTabletInsertionEvent) event).getSourceEvent()) + .registerGeneratedTabletInsertionEvent(); + } + 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 dc2ab1d381fd..b6035c3678e4 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 @@ -118,12 +118,18 @@ 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) { + ((PipeTsFileInsertionEvent) sourceEvent) + .markGeneratedTabletInsertionEventAsTransferred(); + } }); } } @@ -308,7 +314,8 @@ public boolean internallyDecreaseResourceReferenceCount(final String holderMessa protected void eliminateProgressIndex() { if (sourceEvent instanceof PipeTsFileInsertionEvent) { - ((PipeTsFileInsertionEvent) sourceEvent).eliminateProgressIndex(); + final PipeTsFileInsertionEvent tsFileInsertionEvent = (PipeTsFileInsertionEvent) sourceEvent; + tsFileInsertionEvent.eliminateProgressIndex(); } } @@ -391,7 +398,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 ba2a9d275c39..380fc827ee56 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 @@ -109,6 +109,13 @@ 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 AtomicBoolean generatedTabletInsertionEventsParsingCompleted = + new AtomicBoolean(false); + private final AtomicBoolean isTsFileEventCommitted = new AtomicBoolean(false); + private final AtomicBoolean isTransferred = new AtomicBoolean(false); + private final List onTransferredHooks = 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. @@ -295,6 +302,7 @@ private PipeTsFileInsertionEvent( if (shouldReportOnCommit) { eliminateProgressIndex(); } + markAsTransferredOnCommit(); }); } @@ -520,6 +528,75 @@ 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(); + } + + /** Records that one tablet generated from this TsFile has been committed downstream. */ + public void markGeneratedTabletInsertionEventAsTransferred() { + transferredGeneratedTabletInsertionEventCount.incrementAndGet(); + markAsTransferredIfGeneratedTabletEventsCompleted(); + } + + /** Marks the end of tablet generation for this TsFile. */ + public void markGeneratedTabletInsertionEventsParsingCompleted() { + generatedTabletInsertionEventsParsingCompleted.set(true); + markAsTransferredIfGeneratedTabletEventsCompleted(); + } + + private void markAsTransferredIfGeneratedTabletEventsCompleted() { + if ((generatedTabletInsertionEventsParsingCompleted.get() || isTsFileEventCommitted.get()) + && generatedTabletInsertionEventCount.get() > 0 + && transferredGeneratedTabletInsertionEventCount.get() + >= generatedTabletInsertionEventCount.get()) { + markAsTransferred(); + } + } + + private void markAsTransferredOnCommit() { + // The default PipeProcessor implementation iterates over toTabletInsertionEvents() directly, + // so it cannot signal parser completion through consumeTabletInsertionEventsWithRetry(). Its + // source event is committed only after processor execution has returned, which makes this a + // reliable completion boundary for the generated-event counters as well. + isTsFileEventCommitted.set(true); + if (generatedTabletInsertionEventCount.get() == 0) { + markAsTransferred(); + return; + } + markAsTransferredIfGeneratedTabletEventsCompleted(); + } + + private void markAsTransferred() { + final List hooksToRun; + synchronized (onTransferredHooks) { + if (!isTransferred.compareAndSet(false, true)) { + return; + } + hooksToRun = new ArrayList<>(onTransferredHooks); + 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) { + synchronized (onTransferredHooks) { + if (!isTransferred.get()) { + onTransferredHooks.add(hook); + return; + } + } + hook.run(); + } + public PipeTsFileInsertionEvent bindTsFileDedupScopeID(final String tsFileDedupScopeID) { this.tsFileDedupScopeID = tsFileDedupScopeID; return this; @@ -776,6 +853,7 @@ public void consumeTabletInsertionEventsWithRetry( getNextTabletInsertionEventFromSavedProgress(); if (parsedEvent == null) { isTsFileParsingCompleted.set(true); + markGeneratedTabletInsertionEventsParsingCompleted(); releaseTsFileParserMemoryIfReserved(); return; } 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 65d03919ab1b..a68e122a5072 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,40 @@ 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; + ((PipeTsFileInsertionEvent) event.getEvent()) + .addOnTransferredHook(() -> clearTsFileEpochAfterCommit(event.getTsFileEpoch())); } 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 fc4a15a885b5..6637e101f96f 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,598 @@ 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 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 +1185,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 9e0f5a2ecce5..0fb96ff373ee 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 @@ -137,6 +137,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";