From 12bc3f3a5c00265f5c04dc28a4e6c76affe326bf Mon Sep 17 00:00:00 2001 From: Caideyipi <87789683+Caideyipi@users.noreply.github.com> Date: Wed, 22 Jul 2026 18:00:04 +0800 Subject: [PATCH] [Pipe] Add reliable DataNode completion metric --- .../db/pipe/agent/task/PipeDataNodeTask.java | 6 + .../agent/task/PipeDataNodeTaskAgent.java | 101 +++++ .../task/builder/PipeDataNodeTaskBuilder.java | 1 + .../task/connection/PipeEventCollector.java | 23 + .../task/stage/PipeTaskProcessorStage.java | 2 + .../agent/task/stage/PipeTaskSourceStage.java | 15 + .../processor/PipeProcessorSubtask.java | 8 +- .../PipeRealtimePriorityBlockingQueue.java | 4 +- .../common/heartbeat/PipeHeartbeatEvent.java | 94 +++- .../realtime/PipeRealtimeEventFactory.java | 4 + .../PipeDataNodeCompletionOperator.java | 422 ++++++++++++++++++ ...DataNodeRemainingEventAndTimeOperator.java | 4 + .../PipeDataNodeSinglePipeMetrics.java | 128 +++++- .../dataregion/IoTDBDataRegionSource.java | 15 + .../PipeRealtimeDataRegionHybridSource.java | 17 +- .../PipeRealtimeDataRegionLogSource.java | 1 + .../PipeRealtimeDataRegionSource.java | 40 +- .../PipeRealtimeDataRegionTsFileSource.java | 1 + .../assigner/PipeDataRegionAssigner.java | 230 +++++++++- .../PipeInsertionDataNodeListener.java | 77 ++-- .../iotdb/db/storageengine/StorageEngine.java | 17 +- .../storageengine/dataregion/DataRegion.java | 140 ++++++ .../dataregion/memtable/TsFileProcessor.java | 44 +- .../connection/PipeEventCollectorTest.java | 4 +- .../heartbeat/PipeHeartbeatEventTest.java | 83 ++++ .../PipeDataNodeCompletionOperatorTest.java | 252 +++++++++++ .../PipeDataNodeSinglePipeMetricsTest.java | 221 +++++++++ .../PipeRealtimeDataRegionSourceTest.java | 25 ++ .../assigner/PipeDataRegionAssignerTest.java | 75 ++++ .../PipeInsertionDataNodeListenerTest.java | 66 +++ .../dataregion/DataRegionTest.java | 143 ++++++ .../pipe/agent/task/meta/PipeTaskMeta.java | 17 +- .../task/progress/PipeEventCommitManager.java | 10 + .../commons/service/metric/enums/Metric.java | 1 + 34 files changed, 2213 insertions(+), 78 deletions(-) create mode 100644 iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeCompletionOperator.java create mode 100644 iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/common/heartbeat/PipeHeartbeatEventTest.java create mode 100644 iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeCompletionOperatorTest.java create mode 100644 iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeSinglePipeMetricsTest.java create mode 100644 iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/assigner/PipeDataRegionAssignerTest.java create mode 100644 iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/listener/PipeInsertionDataNodeListenerTest.java diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTask.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTask.java index 0d4e50cabef1f..e4731c3426680 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTask.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTask.java @@ -22,6 +22,8 @@ import org.apache.iotdb.commons.pipe.agent.task.PipeTask; import org.apache.iotdb.commons.pipe.agent.task.stage.PipeTaskStage; import org.apache.iotdb.db.i18n.DataNodePipeMessages; +import org.apache.iotdb.db.pipe.agent.task.stage.PipeTaskSourceStage; +import org.apache.iotdb.pipe.api.PipeExtractor; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -109,6 +111,10 @@ public String getPipeName() { return pipeName; } + public PipeExtractor getPipeExtractor() { + return ((PipeTaskSourceStage) sourceStage).getPipeExtractor(); + } + public boolean isCompleted() { return isCompleted; } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java index dce7ee5f6d925..c3c968ecd0e45 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/PipeDataNodeTaskAgent.java @@ -38,11 +38,13 @@ import org.apache.iotdb.commons.pipe.agent.task.meta.PipeMeta; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeRuntimeMeta; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStaticMeta; +import org.apache.iotdb.commons.pipe.agent.task.meta.PipeStatus; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTemporaryMeta; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTemporaryMetaInAgent; import org.apache.iotdb.commons.pipe.agent.task.meta.PipeType; import org.apache.iotdb.commons.pipe.config.PipeConfig; +import org.apache.iotdb.commons.pipe.config.constant.PipeProcessorConstant; import org.apache.iotdb.commons.pipe.config.constant.PipeSinkConstant; import org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant; import org.apache.iotdb.commons.pipe.resource.log.PipeLogger; @@ -61,7 +63,10 @@ import org.apache.iotdb.db.pipe.resource.memory.PipeMemoryManager; import org.apache.iotdb.db.pipe.resource.tsfile.PipeTsFileResourceManager; import org.apache.iotdb.db.pipe.source.dataregion.DataRegionListeningFilter; +import org.apache.iotdb.db.pipe.source.dataregion.IoTDBDataRegionSource; +import org.apache.iotdb.db.pipe.source.dataregion.realtime.PipeRealtimeDataRegionSource; import org.apache.iotdb.db.pipe.source.dataregion.realtime.listener.PipeInsertionDataNodeListener; +import org.apache.iotdb.db.pipe.source.schemaregion.IoTDBSchemaRegionSource; import org.apache.iotdb.db.pipe.source.schemaregion.SchemaRegionListeningFilter; import org.apache.iotdb.db.protocol.client.ConfigNodeClient; import org.apache.iotdb.db.protocol.client.ConfigNodeClientManager; @@ -73,6 +78,7 @@ import org.apache.iotdb.mpp.rpc.thrift.TDataNodeHeartbeatResp; import org.apache.iotdb.mpp.rpc.thrift.TPipeHeartbeatReq; import org.apache.iotdb.mpp.rpc.thrift.TPushPipeMetaRespExceptionMessage; +import org.apache.iotdb.pipe.api.PipeExtractor; import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameters; import org.apache.iotdb.pipe.api.exception.PipeException; import org.apache.iotdb.rpc.TSStatusCode; @@ -115,6 +121,20 @@ public class PipeDataNodeTaskAgent extends PipeTaskAgent { private static final Logger LOGGER = LoggerFactory.getLogger(PipeDataNodeTaskAgent.class); + private static final Set COMPLETION_SUPPORTED_SOURCES = + Set.of( + BuiltinPipePlugin.IOTDB_EXTRACTOR.getPipePluginName(), + BuiltinPipePlugin.IOTDB_SOURCE.getPipePluginName()); + private static final Set COMPLETION_SUPPORTED_SINKS = + Set.of( + BuiltinPipePlugin.IOTDB_THRIFT_CONNECTOR.getPipePluginName(), + BuiltinPipePlugin.IOTDB_THRIFT_SSL_CONNECTOR.getPipePluginName(), + BuiltinPipePlugin.IOTDB_THRIFT_SYNC_CONNECTOR.getPipePluginName(), + BuiltinPipePlugin.IOTDB_THRIFT_ASYNC_CONNECTOR.getPipePluginName(), + BuiltinPipePlugin.IOTDB_THRIFT_SINK.getPipePluginName(), + BuiltinPipePlugin.IOTDB_THRIFT_SSL_SINK.getPipePluginName(), + BuiltinPipePlugin.IOTDB_THRIFT_SYNC_SINK.getPipePluginName(), + BuiltinPipePlugin.IOTDB_THRIFT_ASYNC_SINK.getPipePluginName()); protected static final IoTDBConfig CONFIG = IoTDBDescriptor.getInstance().getConfig(); @@ -570,6 +590,87 @@ public Set getPipeTaskRegionIdSet(final String pipeName, final long cre : pipeMeta.getRuntimeMeta().getConsensusGroupId2TaskMetaMap().keySet(); } + public Pair> getPipeCompletionSnapshot( + final String pipeName, final long creationTime) { + if (!tryReadLockWithTimeOutInMs(10)) { + return new Pair<>(false, Collections.emptyMap()); + } + + try { + final PipeMeta pipeMeta = pipeMetaKeeper.getPipeMeta(pipeName, creationTime); + if (pipeMeta == null) { + return new Pair<>(false, Collections.emptyMap()); + } + + final PipeStaticMeta staticMeta = pipeMeta.getStaticMeta(); + final String sourceName = + staticMeta + .getSourceParameters() + .getStringOrDefault( + Arrays.asList(PipeSourceConstant.EXTRACTOR_KEY, PipeSourceConstant.SOURCE_KEY), + BuiltinPipePlugin.IOTDB_EXTRACTOR.getPipePluginName()); + final String processorName = + staticMeta + .getProcessorParameters() + .getStringOrDefault( + PipeProcessorConstant.PROCESSOR_KEY, + BuiltinPipePlugin.DO_NOTHING_PROCESSOR.getPipePluginName()); + final PipeParameters sinkParameters = staticMeta.getSinkParameters(); + final String sinkName = PipeSinkConstant.getConnectorOrSinkNameWithDefault(sinkParameters); + boolean supported = + staticMeta.getPipeType() == PipeType.USER + && COMPLETION_SUPPORTED_SOURCES.stream() + .anyMatch(name -> name.equalsIgnoreCase(sourceName)) + && BuiltinPipePlugin.DO_NOTHING_PROCESSOR + .getPipePluginName() + .equalsIgnoreCase(processorName) + && COMPLETION_SUPPORTED_SINKS.stream() + .anyMatch(name -> name.equalsIgnoreCase(sinkName)) + && PipeSinkConstant.CONNECTOR_LOAD_TSFILE_STRATEGY_SYNC_VALUE.equalsIgnoreCase( + sinkParameters.getStringOrDefault( + Arrays.asList( + PipeSinkConstant.CONNECTOR_LOAD_TSFILE_STRATEGY_KEY, + PipeSinkConstant.SINK_LOAD_TSFILE_STRATEGY_KEY), + PipeSinkConstant.CONNECTOR_LOAD_TSFILE_STRATEGY_SYNC_VALUE)) + && pipeMeta.getRuntimeMeta().getStatus().get() == PipeStatus.RUNNING + && pipeMeta.getRuntimeMeta().getNodeId2PipeRuntimeExceptionMap().isEmpty() + && pipeMeta.getRuntimeMeta().getConsensusGroupId2TaskMetaMap().values().stream() + .noneMatch(PipeTaskMeta::hasExceptionMessages); + + final Map dataRegionId2Source = new HashMap<>(); + final Map pipeTaskMap = pipeTaskManager.getPipeTasks(staticMeta); + if (pipeTaskMap != null) { + for (final Map.Entry entry : pipeTaskMap.entrySet()) { + if (!(entry.getValue() instanceof PipeDataNodeTask)) { + supported = false; + continue; + } + final PipeExtractor extractor = ((PipeDataNodeTask) entry.getValue()).getPipeExtractor(); + if (!(extractor instanceof IoTDBDataRegionSource)) { + if (!(extractor instanceof IoTDBSchemaRegionSource)) { + supported = false; + } + continue; + } + final IoTDBDataRegionSource dataRegionSource = (IoTDBDataRegionSource) extractor; + final PipeRealtimeDataRegionSource realtimeSource = + dataRegionSource.getRealtimeSourceForCompletion(); + if (realtimeSource == null) { + supported = false; + } else { + dataRegionId2Source.put(entry.getKey(), realtimeSource); + } + if (!dataRegionSource.isReadyForCompletion()) { + supported = false; + } + } + } + return new Pair<>(supported, dataRegionId2Source); + } finally { + releaseReadLock(); + } + } + public boolean hasPipeReleaseRegionRelatedResource(final int consensusGroupId) { if (!tryReadLockWithTimeOut(10)) { LOGGER.warn(DataNodePipeMessages.FAILED_TO_CHECK_IF_PIPE_HAS_RELEASE, consensusGroupId); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/builder/PipeDataNodeTaskBuilder.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/builder/PipeDataNodeTaskBuilder.java index 98157a2a49a09..01efe08fdd934 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/builder/PipeDataNodeTaskBuilder.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/builder/PipeDataNodeTaskBuilder.java @@ -125,6 +125,7 @@ public PipeDataNodeTask build() { pipeStaticMeta.getCreationTime(), blendUserAndSystemParameters(pipeStaticMeta.getProcessorParameters(), pipeTaskMeta), regionId, + sourceStage.getCompletionSourceId(), sourceStage.getEventSupplier(), sinkStage.getPipeSinkPendingQueue(), PROCESSOR_EXECUTOR, 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 5ecd81d92b412..457ffd749507a 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 @@ -34,6 +34,7 @@ import org.apache.iotdb.db.pipe.event.common.tablet.PipeRawTabletInsertionEvent; import org.apache.iotdb.db.pipe.event.common.terminate.PipeTerminateEvent; import org.apache.iotdb.db.pipe.event.common.tsfile.PipeTsFileInsertionEvent; +import org.apache.iotdb.db.pipe.metric.overview.PipeDataNodeSinglePipeMetrics; import org.apache.iotdb.db.pipe.source.schemaregion.IoTDBSchemaRegionSource; import org.apache.iotdb.db.pipe.source.schemaregion.PipePlanTablePrivilegeParseVisitor; import org.apache.iotdb.db.pipe.source.schemaregion.PipePlanTreePrivilegeParseVisitor; @@ -58,6 +59,8 @@ public class PipeEventCollector implements EventCollector { private final int regionId; + private final long completionSourceId; + private final boolean forceTabletFormat; private final boolean skipParsing; @@ -72,12 +75,14 @@ public PipeEventCollector( final UnboundedBlockingPendingQueue pendingQueue, final long creationTime, final int regionId, + final long completionSourceId, final boolean forceTabletFormat, final boolean skipParsing, final boolean isUsedInConsensusPipe) { this.pendingQueue = pendingQueue; this.creationTime = creationTime; this.regionId = regionId; + this.completionSourceId = completionSourceId; this.forceTabletFormat = forceTabletFormat; this.skipParsing = skipParsing; this.isUsedForConsensusPipe = isUsedInConsensusPipe; @@ -236,6 +241,7 @@ private void collectEvent(final Event event) { if (event instanceof EnrichedEvent) { final EnrichedEvent enrichedEvent = (EnrichedEvent) event; if (!enrichedEvent.increaseReferenceCount(PipeEventCollector.class.getName())) { + markDataRegionCompletionInvalid(enrichedEvent); LOGGER.warn( DataNodePipeMessages.PIPEEVENTCOLLECTOR_THE_EVENT_IS_ALREADY_RELEASED_SKIPPING, event); isFailedToIncreaseReferenceCount = true; @@ -245,6 +251,9 @@ private void collectEvent(final Event event) { // Assign a commit id for this event in order to report progress in order. PipeEventCommitManager.getInstance() .enrichWithCommitterKeyAndCommitId(enrichedEvent, creationTime, regionId); + if (enrichedEvent.needToCommit() && enrichedEvent.getCommitterKey() == null) { + markDataRegionCompletionInvalid(enrichedEvent); + } // Assign a rebootTime for iotConsensusV2 enrichedEvent.setRebootTimes(PipeDataNodeAgent.runtime().getRebootTimes()); @@ -265,6 +274,20 @@ private void collectEvent(final Event event) { if (pendingQueue.offer(event)) { collectInvocationCount.incrementAndGet(); + } else if (event instanceof EnrichedEvent) { + markDataRegionCompletionInvalid((EnrichedEvent) event); + } + } + + private void markDataRegionCompletionInvalid(final EnrichedEvent event) { + if (!(event instanceof PipeHeartbeatEvent) && event.getPipeName() != null && regionId >= 0) { + PipeDataNodeSinglePipeMetrics.getInstance() + .markDataRegionInvalid( + event.getPipeName(), + creationTime, + regionId, + event.getPipeTaskMeta(), + completionSourceId); } } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/stage/PipeTaskProcessorStage.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/stage/PipeTaskProcessorStage.java index 85a3ff4cd3c60..a84ff51421830 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/stage/PipeTaskProcessorStage.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/stage/PipeTaskProcessorStage.java @@ -67,6 +67,7 @@ public PipeTaskProcessorStage( final long creationTime, final PipeParameters pipeProcessorParameters, final int regionId, + final long completionSourceId, final EventSupplier pipeSourceInputEventSupplier, final UnboundedBlockingPendingQueue pipeSinkOutputPendingQueue, final PipeProcessorSubtaskExecutor executor, @@ -111,6 +112,7 @@ public PipeTaskProcessorStage( pipeSinkOutputPendingQueue, creationTime, regionId, + completionSourceId, forceTabletFormat, skipParsing, isUsedForConsensusPipe); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/stage/PipeTaskSourceStage.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/stage/PipeTaskSourceStage.java index 460fad472d12c..3dcb5ee4c5c3b 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/stage/PipeTaskSourceStage.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/stage/PipeTaskSourceStage.java @@ -28,6 +28,8 @@ import org.apache.iotdb.commons.pipe.config.plugin.env.PipeTaskSourceRuntimeEnvironment; import org.apache.iotdb.db.i18n.DataNodePipeMessages; import org.apache.iotdb.db.pipe.agent.PipeDataNodeAgent; +import org.apache.iotdb.db.pipe.source.dataregion.IoTDBDataRegionSource; +import org.apache.iotdb.db.pipe.source.dataregion.realtime.PipeRealtimeDataRegionSource; import org.apache.iotdb.db.storageengine.StorageEngine; import org.apache.iotdb.pipe.api.PipeExtractor; import org.apache.iotdb.pipe.api.customizer.parameter.PipeParameterValidator; @@ -107,4 +109,17 @@ public void dropSubtask() throws PipeException { public EventSupplier getEventSupplier() { return pipeExtractor::supply; } + + public PipeExtractor getPipeExtractor() { + return pipeExtractor; + } + + public long getCompletionSourceId() { + if (!(pipeExtractor instanceof IoTDBDataRegionSource)) { + return Long.MIN_VALUE; + } + final PipeRealtimeDataRegionSource realtimeSource = + ((IoTDBDataRegionSource) pipeExtractor).getRealtimeSourceForCompletion(); + return realtimeSource == null ? Long.MIN_VALUE : realtimeSource.getCompletionSourceId(); + } } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/processor/PipeProcessorSubtask.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/processor/PipeProcessorSubtask.java index 8c26726bfb204..ccc083bff1ec6 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/processor/PipeProcessorSubtask.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/processor/PipeProcessorSubtask.java @@ -204,7 +204,13 @@ protected boolean executeOnce() throws Exception { .markTsFileCollectInvocationCount( pipeNameWithCreationTime, outputEventCollector.getCollectInvocationCount()); } else if (event instanceof PipeHeartbeatEvent) { - pipeProcessor.process(event, outputEventCollector); + // A completion barrier is an internal ordering event. It must not be swallowed or + // transformed by a user processor. + if (((PipeHeartbeatEvent) event).isCompletionBarrier()) { + outputEventCollector.collect(event); + } else { + pipeProcessor.process(event, outputEventCollector); + } ((PipeHeartbeatEvent) event).onProcessed(); PipeProcessorMetrics.getInstance().markPipeHeartbeatEvent(taskID); } else { diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeRealtimePriorityBlockingQueue.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeRealtimePriorityBlockingQueue.java index 8641dbc78674e..f224f09425d43 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeRealtimePriorityBlockingQueue.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/agent/task/subtask/sink/PipeRealtimePriorityBlockingQueue.java @@ -83,7 +83,9 @@ public boolean offer(final Event event) { return true; } - if (event instanceof PipeHeartbeatEvent && super.peekLast() instanceof PipeHeartbeatEvent) { + if (event instanceof PipeHeartbeatEvent + && !((PipeHeartbeatEvent) event).isCompletionBarrier() + && super.peekLast() instanceof PipeHeartbeatEvent) { // We can NOT keep too many PipeHeartbeatEvent in bufferQueue because they may cause OOM. ((EnrichedEvent) event).decreaseReferenceCount(PipeEventCollector.class.getName(), false); return false; diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/heartbeat/PipeHeartbeatEvent.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/heartbeat/PipeHeartbeatEvent.java index 33697f6d5ab85..cbe264b2ca594 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/heartbeat/PipeHeartbeatEvent.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/common/heartbeat/PipeHeartbeatEvent.java @@ -42,6 +42,10 @@ public class PipeHeartbeatEvent extends EnrichedEvent { private static final Logger LOGGER = LoggerFactory.getLogger(PipeHeartbeatEvent.class); private final int dataRegionId; + private final boolean completionBarrier; + private long assignerEpoch = Long.MIN_VALUE; + private long dataGeneration = Long.MIN_VALUE; + private long completionSourceId = Long.MIN_VALUE; private long timePublished; private long timeAssigned; @@ -63,9 +67,15 @@ public class PipeHeartbeatEvent extends EnrichedEvent { private final boolean shouldPrintMessage; public PipeHeartbeatEvent(final int dataRegionId, final boolean shouldPrintMessage) { + this(dataRegionId, shouldPrintMessage, false); + } + + public PipeHeartbeatEvent( + final int dataRegionId, final boolean shouldPrintMessage, final boolean completionBarrier) { super(null, 0, null, null, null, null, null, null, true, Long.MIN_VALUE, Long.MAX_VALUE); this.dataRegionId = dataRegionId; this.shouldPrintMessage = shouldPrintMessage; + this.completionBarrier = completionBarrier; } public PipeHeartbeatEvent( @@ -74,7 +84,11 @@ public PipeHeartbeatEvent( final PipeTaskMeta pipeTaskMeta, final int dataRegionId, final long timePublished, - final boolean shouldPrintMessage) { + final boolean shouldPrintMessage, + final boolean completionBarrier, + final long assignerEpoch, + final long dataGeneration, + final long completionSourceId) { super( pipeName, creationTime, @@ -90,6 +104,13 @@ public PipeHeartbeatEvent( this.dataRegionId = dataRegionId; this.timePublished = timePublished; this.shouldPrintMessage = shouldPrintMessage; + this.completionBarrier = completionBarrier; + this.assignerEpoch = assignerEpoch; + this.dataGeneration = dataGeneration; + this.completionSourceId = completionSourceId; + if (completionBarrier && completionSourceId != Long.MIN_VALUE) { + addCompletionOnCommittedHook(); + } } @Override @@ -136,7 +157,16 @@ public EnrichedEvent shallowCopySelfAndBindPipeTaskMetaForProgressReport( // Should record PipeTaskMeta, for sometimes HeartbeatEvents should report exceptions. // Here we ignore parameters `pattern`, `startTime`, and `endTime`. return new PipeHeartbeatEvent( - pipeName, creationTime, pipeTaskMeta, dataRegionId, timePublished, shouldPrintMessage); + pipeName, + creationTime, + pipeTaskMeta, + dataRegionId, + timePublished, + shouldPrintMessage, + completionBarrier, + assignerEpoch, + dataGeneration, + completionSourceId); } @Override @@ -160,6 +190,54 @@ public boolean isShouldPrintMessage() { return shouldPrintMessage; } + public boolean isCompletionBarrier() { + return completionBarrier; + } + + public void bindCompletionBarrier(final long assignerEpoch, final long dataGeneration) { + if (completionBarrier) { + this.assignerEpoch = assignerEpoch; + this.dataGeneration = dataGeneration; + } + } + + public void bindCompletionSource(final long completionSourceId) { + if (completionBarrier + && pipeName != null + && this.completionSourceId == Long.MIN_VALUE + && completionSourceId != Long.MIN_VALUE) { + this.completionSourceId = completionSourceId; + addCompletionOnCommittedHook(); + } + } + + private void addCompletionOnCommittedHook() { + addOnCommittedHook( + () -> + PipeDataNodeSinglePipeMetrics.getInstance() + .markDataRegionCompleted( + pipeName, + creationTime, + dataRegionId, + pipeTaskMeta, + assignerEpoch, + dataGeneration, + completionSourceId, + getCommitterKey())); + } + + public long getAssignerEpoch() { + return assignerEpoch; + } + + public long getDataGeneration() { + return dataGeneration; + } + + public long getCompletionSourceId() { + return completionSourceId; + } + /////////////////////////////// Delay Reporting /////////////////////////////// public void onPublished() { @@ -224,10 +302,10 @@ public void recordConnectorQueueSize(final UnboundedBlockingPendingQueue /////////////////////////////// For Commit Ordering /////////////////////////////// - /** {@link PipeHeartbeatEvent}s do not need to be committed in order. */ + /** Only completion barriers need ordered commit. Periodic heartbeats remain best-effort. */ @Override public boolean needToCommit() { - return false; + return completionBarrier && Objects.nonNull(pipeName) && completionSourceId != Long.MIN_VALUE; } /////////////////////////////// Object /////////////////////////////// @@ -278,6 +356,14 @@ public String toString() { + pipeName + "', dataRegionId=" + dataRegionId + + ", completionBarrier=" + + completionBarrier + + ", assignerEpoch=" + + assignerEpoch + + ", dataGeneration=" + + dataGeneration + + ", completionSourceId=" + + completionSourceId + ", startTime=" + startTimeMessage + ", publishedToAssigned=" diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/realtime/PipeRealtimeEventFactory.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/realtime/PipeRealtimeEventFactory.java index 540589326ffec..6784d5a654188 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/realtime/PipeRealtimeEventFactory.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/event/realtime/PipeRealtimeEventFactory.java @@ -62,6 +62,10 @@ public static PipeRealtimeEvent createRealtimeEvent( new PipeHeartbeatEvent(dataRegionId, shouldPrintMessage), null, null); } + public static PipeRealtimeEvent createCompletionBarrierEvent(final int dataRegionId) { + return new PipeRealtimeEvent(new PipeHeartbeatEvent(dataRegionId, false, true), null, null); + } + public static PipeRealtimeEvent createRealtimeEvent(final AbstractDeleteDataNode node) { PipeDeleteDataNodeEvent deleteDataNodeEvent = new PipeDeleteDataNodeEvent(node, node.isGeneratedByPipe()); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeCompletionOperator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeCompletionOperator.java new file mode 100644 index 0000000000000..87084efaee657 --- /dev/null +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeCompletionOperator.java @@ -0,0 +1,422 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.iotdb.db.pipe.metric.overview; + +import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta; +import org.apache.iotdb.commons.pipe.agent.task.progress.CommitterKey; +import org.apache.iotdb.commons.pipe.agent.task.progress.PipeEventCommitManager; +import org.apache.iotdb.db.pipe.source.dataregion.realtime.PipeRealtimeDataRegionSource; +import org.apache.iotdb.db.pipe.source.dataregion.realtime.assigner.PipeDataRegionAssigner; + +import org.apache.tsfile.utils.Pair; + +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicReference; +import java.util.concurrent.locks.ReentrantReadWriteLock; +import java.util.function.BooleanSupplier; +import java.util.function.Predicate; +import java.util.function.Supplier; + +/** + * Calculates whether all local DataRegion tasks of a pipe have committed their latest full-flush + * completion barriers. + * + *

The operator deliberately fails closed. It takes two task-topology snapshots and two + * membership/state snapshots. Any change in task identity, source identity, health, remaining + * events, degradation, publication generation, publication failures, or committer incarnation makes + * the result incomplete. + */ +class PipeDataNodeCompletionOperator { + + private final BooleanSupplier isRemainingEventCountZero; + private final Supplier>> + taskSnapshotSupplier; + private final Predicate isCurrentCommitterKey; + + private final ReentrantReadWriteLock membershipLock = new ReentrantReadWriteLock(); + private final Map registeredSourceMap = new HashMap<>(); + private final Map dataRegionId2State = new HashMap<>(); + private long membershipVersion = 0; + + PipeDataNodeCompletionOperator( + final BooleanSupplier isRemainingEventCountZero, + final Supplier>> + taskSnapshotSupplier) { + this( + isRemainingEventCountZero, + taskSnapshotSupplier, + PipeEventCommitManager.getInstance()::isCurrentCommitterKey); + } + + PipeDataNodeCompletionOperator( + final BooleanSupplier isRemainingEventCountZero, + final Supplier>> + taskSnapshotSupplier, + final Predicate isCurrentCommitterKey) { + this.isRemainingEventCountZero = isRemainingEventCountZero; + this.taskSnapshotSupplier = taskSnapshotSupplier; + this.isCurrentCommitterKey = isCurrentCommitterKey; + } + + void registerDataRegionSource(final PipeRealtimeDataRegionSource source) { + membershipLock.writeLock().lock(); + try { + registeredSourceMap.put(source.getDataRegionId(), source); + dataRegionId2State.computeIfPresent( + source.getDataRegionId(), (regionId, state) -> state.source == source ? state : null); + membershipVersion++; + } finally { + membershipLock.writeLock().unlock(); + } + } + + void deregisterDataRegionSource(final PipeRealtimeDataRegionSource source) { + membershipLock.writeLock().lock(); + try { + if (registeredSourceMap.get(source.getDataRegionId()) == source) { + registeredSourceMap.remove(source.getDataRegionId()); + dataRegionId2State.computeIfPresent( + source.getDataRegionId(), (regionId, state) -> state.source == source ? null : state); + membershipVersion++; + } + } finally { + membershipLock.writeLock().unlock(); + } + } + + void register(final PipeRealtimeDataRegionSource source, final PipeDataRegionAssigner assigner) { + membershipLock.writeLock().lock(); + try { + if (registeredSourceMap.get(source.getDataRegionId()) == source) { + dataRegionId2State.put( + source.getDataRegionId(), + new DataRegionCompletionState(source, assigner, isCurrentCommitterKey)); + membershipVersion++; + } + } finally { + membershipLock.writeLock().unlock(); + } + } + + void deregister( + final PipeRealtimeDataRegionSource source, final PipeDataRegionAssigner assigner) { + membershipLock.writeLock().lock(); + try { + final DataRegionCompletionState state = dataRegionId2State.get(source.getDataRegionId()); + if (state != null && state.matches(source, assigner)) { + dataRegionId2State.remove(source.getDataRegionId()); + membershipVersion++; + } + } finally { + membershipLock.writeLock().unlock(); + } + } + + void markCompleted( + final int dataRegionId, + final PipeTaskMeta pipeTaskMeta, + final long assignerEpoch, + final long generation, + final long completionSourceId, + final CommitterKey committerKey) { + membershipLock.readLock().lock(); + try { + final DataRegionCompletionState state = dataRegionId2State.get(dataRegionId); + if (state != null) { + state.markCompleted( + pipeTaskMeta, assignerEpoch, generation, completionSourceId, committerKey); + } + } finally { + membershipLock.readLock().unlock(); + } + } + + void markInvalid( + final int dataRegionId, final PipeTaskMeta pipeTaskMeta, final long completionSourceId) { + membershipLock.readLock().lock(); + try { + final DataRegionCompletionState state = dataRegionId2State.get(dataRegionId); + if (state != null) { + state.markInvalid(pipeTaskMeta, completionSourceId); + } + } finally { + membershipLock.readLock().unlock(); + } + } + + public long getCompletion() { + try { + if (!isRemainingEventCountZero.getAsBoolean()) { + return 0; + } + + final Pair> firstTaskSnapshot = + taskSnapshotSupplier.get(); + if (!isTaskSnapshotSupported(firstTaskSnapshot)) { + return 0; + } + final Map expectedSources = + new HashMap<>(firstTaskSnapshot.getRight()); + + final MembershipObservation membershipObservation = observeMembership(expectedSources); + if (membershipObservation == null) { + return 0; + } + + final Pair> secondTaskSnapshot = + taskSnapshotSupplier.get(); + if (!isTaskSnapshotSupported(secondTaskSnapshot) + || !expectedSources.equals(secondTaskSnapshot.getRight()) + || !isRemainingEventCountZero.getAsBoolean()) { + return 0; + } + + return verifyMembership(expectedSources, membershipObservation) ? 1 : 0; + } catch (final RuntimeException e) { + // A metric must fail closed when task topology or lifecycle is concurrently changing. + return 0; + } + } + + private MembershipObservation observeMembership( + final Map expectedSources) { + membershipLock.readLock().lock(); + try { + if (!registeredSourceMap.equals(expectedSources) + || !dataRegionId2State.keySet().equals(expectedSources.keySet())) { + return null; + } + + final Map stateObservations = new HashMap<>(); + for (final Map.Entry entry : + expectedSources.entrySet()) { + final DataRegionCompletionState state = dataRegionId2State.get(entry.getKey()); + final StateObservation observation = + state == null ? null : state.observeCompletion(entry.getValue()); + if (observation == null) { + return null; + } + stateObservations.put(entry.getKey(), observation); + } + return new MembershipObservation(membershipVersion, stateObservations); + } finally { + membershipLock.readLock().unlock(); + } + } + + private boolean verifyMembership( + final Map expectedSources, + final MembershipObservation membershipObservation) { + membershipLock.readLock().lock(); + try { + if (membershipVersion != membershipObservation.membershipVersion + || !registeredSourceMap.equals(expectedSources) + || !dataRegionId2State.keySet().equals(expectedSources.keySet())) { + return false; + } + + for (final Map.Entry entry : + expectedSources.entrySet()) { + final DataRegionCompletionState state = dataRegionId2State.get(entry.getKey()); + if (state == null + || !state.isStillCompleted( + entry.getValue(), membershipObservation.stateObservations.get(entry.getKey()))) { + return false; + } + } + return true; + } finally { + membershipLock.readLock().unlock(); + } + } + + private static boolean isTaskSnapshotSupported( + final Pair> snapshot) { + return snapshot != null + && Boolean.TRUE.equals(snapshot.getLeft()) + && snapshot.getRight() != null; + } + + private static class DataRegionCompletionState { + + private final PipeRealtimeDataRegionSource source; + private final PipeDataRegionAssigner assigner; + private final PipeTaskMeta pipeTaskMeta; + private final long completionSourceId; + private final long assignerEpoch; + private final long initialPublicationFailureEpoch; + private final Predicate isCurrentCommitterKey; + private final AtomicBoolean valid = new AtomicBoolean(true); + private final AtomicReference committedBarrier = new AtomicReference<>(); + + private DataRegionCompletionState( + final PipeRealtimeDataRegionSource source, + final PipeDataRegionAssigner assigner, + final Predicate isCurrentCommitterKey) { + this.source = source; + this.assigner = assigner; + pipeTaskMeta = source.getPipeTaskMeta(); + completionSourceId = source.getCompletionSourceId(); + assignerEpoch = assigner.getAssignerEpoch(); + initialPublicationFailureEpoch = assigner.getPublicationFailureEpoch(); + this.isCurrentCommitterKey = isCurrentCommitterKey; + } + + private boolean matches( + final PipeRealtimeDataRegionSource source, final PipeDataRegionAssigner assigner) { + return this.source == source && this.assigner == assigner; + } + + private void markCompleted( + final PipeTaskMeta pipeTaskMeta, + final long assignerEpoch, + final long generation, + final long completionSourceId, + final CommitterKey committerKey) { + if (this.pipeTaskMeta != pipeTaskMeta + || this.assignerEpoch != assignerEpoch + || this.completionSourceId != completionSourceId + || committerKey == null) { + return; + } + + while (true) { + final CommittedBarrier previous = committedBarrier.get(); + if (previous != null && previous.generation > generation) { + return; + } + final CommittedBarrier next = new CommittedBarrier(generation, committerKey); + if (committedBarrier.compareAndSet(previous, next)) { + return; + } + } + } + + private void markInvalid(final PipeTaskMeta pipeTaskMeta, final long completionSourceId) { + if (this.pipeTaskMeta == pipeTaskMeta && this.completionSourceId == completionSourceId) { + valid.set(false); + } + } + + private StateObservation observeCompletion(final PipeRealtimeDataRegionSource expectedSource) { + final long sourceStateVersion = source.getCompletionStateVersion(); + final long exceptionMessageVersion = + pipeTaskMeta == null ? Long.MIN_VALUE : pipeTaskMeta.getExceptionMessageVersion(); + final long publicationFailureEpoch = assigner.getPublicationFailureEpoch(); + final long publishedGeneration = assigner.getPublishedDataGeneration(); + final CommittedBarrier barrier = committedBarrier.get(); + + if (expectedSource != source + || !valid.get() + || source.getPipeTaskMeta() != pipeTaskMeta + || source.getCompletionSourceId() != completionSourceId + || pipeTaskMeta == null + || pipeTaskMeta.hasExceptionMessages() + || source.isTsFileEpochDegraded() + || assigner.getAssignerEpoch() != assignerEpoch + || publicationFailureEpoch != initialPublicationFailureEpoch + || barrier == null + || barrier.generation < publishedGeneration + || !isCurrentCommitterKey.test(barrier.committerKey)) { + return null; + } + + return new StateObservation( + this, + barrier, + sourceStateVersion, + exceptionMessageVersion, + publicationFailureEpoch, + publishedGeneration); + } + + private boolean isStillCompleted( + final PipeRealtimeDataRegionSource expectedSource, final StateObservation observation) { + return observation != null + && observation.state == this + && expectedSource == source + && valid.get() + && source.getPipeTaskMeta() == pipeTaskMeta + && source.getCompletionSourceId() == completionSourceId + && pipeTaskMeta != null + && !pipeTaskMeta.hasExceptionMessages() + && pipeTaskMeta.getExceptionMessageVersion() == observation.exceptionMessageVersion + && !source.isTsFileEpochDegraded() + && source.getCompletionStateVersion() == observation.sourceStateVersion + && assigner.getAssignerEpoch() == assignerEpoch + && assigner.getPublicationFailureEpoch() == initialPublicationFailureEpoch + && assigner.getPublicationFailureEpoch() == observation.publicationFailureEpoch + && assigner.getPublishedDataGeneration() == observation.publishedGeneration + && committedBarrier.get() == observation.barrier + && observation.barrier.generation >= observation.publishedGeneration + && isCurrentCommitterKey.test(observation.barrier.committerKey); + } + } + + private static class MembershipObservation { + + private final long membershipVersion; + private final Map stateObservations; + + private MembershipObservation( + final long membershipVersion, final Map stateObservations) { + this.membershipVersion = membershipVersion; + this.stateObservations = stateObservations; + } + } + + private static class StateObservation { + + private final DataRegionCompletionState state; + private final CommittedBarrier barrier; + private final long sourceStateVersion; + private final long exceptionMessageVersion; + private final long publicationFailureEpoch; + private final long publishedGeneration; + + private StateObservation( + final DataRegionCompletionState state, + final CommittedBarrier barrier, + final long sourceStateVersion, + final long exceptionMessageVersion, + final long publicationFailureEpoch, + final long publishedGeneration) { + this.state = state; + this.barrier = barrier; + this.sourceStateVersion = sourceStateVersion; + this.exceptionMessageVersion = exceptionMessageVersion; + this.publicationFailureEpoch = publicationFailureEpoch; + this.publishedGeneration = publishedGeneration; + } + } + + private static class CommittedBarrier { + + private final long generation; + private final CommitterKey committerKey; + + private CommittedBarrier(final long generation, final CommitterKey committerKey) { + this.generation = generation; + this.committerKey = committerKey; + } + } +} diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeRemainingEventAndTimeOperator.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeRemainingEventAndTimeOperator.java index a79be559dd86d..30cfcdba7b536 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeRemainingEventAndTimeOperator.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeRemainingEventAndTimeOperator.java @@ -198,6 +198,10 @@ void register(final IoTDBDataRegionSource source) { dataRegionSources.add(source); } + void deregister(final IoTDBDataRegionSource source) { + dataRegionSources.remove(source); + } + void register(final IoTDBSchemaRegionSource source) { schemaRegionSources.add(source); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeSinglePipeMetrics.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeSinglePipeMetrics.java index b6efd64ed2c3c..b779f6492588e 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeSinglePipeMetrics.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeSinglePipeMetrics.java @@ -19,6 +19,8 @@ package org.apache.iotdb.db.pipe.metric.overview; +import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta; +import org.apache.iotdb.commons.pipe.agent.task.progress.CommitterKey; import org.apache.iotdb.commons.pipe.agent.task.progress.PipeEventCommitManager; import org.apache.iotdb.commons.service.metric.enums.Metric; import org.apache.iotdb.commons.service.metric.enums.Tag; @@ -27,6 +29,8 @@ import org.apache.iotdb.db.pipe.resource.PipeDataNodeResourceManager; import org.apache.iotdb.db.pipe.resource.tsfile.PipeTsFileResourceManager; import org.apache.iotdb.db.pipe.source.dataregion.IoTDBDataRegionSource; +import org.apache.iotdb.db.pipe.source.dataregion.realtime.PipeRealtimeDataRegionSource; +import org.apache.iotdb.db.pipe.source.dataregion.realtime.assigner.PipeDataRegionAssigner; import org.apache.iotdb.db.pipe.source.schemaregion.IoTDBSchemaRegionSource; import org.apache.iotdb.metrics.AbstractMetricService; import org.apache.iotdb.metrics.metricsets.IMetricSet; @@ -52,6 +56,8 @@ public class PipeDataNodeSinglePipeMetrics implements IMetricSet { public final Map remainingEventAndTimeOperatorMap = new ConcurrentHashMap<>(); + final Map completionOperatorMap = + new ConcurrentHashMap<>(); //////////////////////////// bindTo & unbindFrom (metric framework) //////////////////////////// @@ -90,6 +96,19 @@ private void createAutoGauge(final String pipeID) { Tag.CREATION_TIME.toString(), String.valueOf(operator.getCreationTime())); + final PipeDataNodeCompletionOperator completionOperator = completionOperatorMap.get(pipeID); + if (completionOperator != null) { + metricService.createAutoGauge( + Metric.PIPE_DATANODE_COMPLETION_READY.toString(), + MetricLevel.IMPORTANT, + completionOperator, + PipeDataNodeCompletionOperator::getCompletion, + Tag.NAME.toString(), + operator.getPipeName(), + Tag.CREATION_TIME.toString(), + String.valueOf(operator.getCreationTime())); + } + // Resources metricService.createAutoGauge( Metric.PIPE_FLOATING_MEMORY_USAGE.toString(), @@ -163,6 +182,15 @@ private void removeAutoGauge(final String pipeID) { operator.getPipeName(), Tag.CREATION_TIME.toString(), String.valueOf(operator.getCreationTime())); + if (completionOperatorMap.containsKey(pipeID)) { + metricService.remove( + MetricType.AUTO_GAUGE, + Metric.PIPE_DATANODE_COMPLETION_READY.toString(), + Tag.NAME.toString(), + operator.getPipeName(), + Tag.CREATION_TIME.toString(), + String.valueOf(operator.getCreationTime())); + } metricService.remove( MetricType.AUTO_GAUGE, Metric.PIPE_FLOATING_MEMORY_USAGE.toString(), @@ -194,7 +222,6 @@ private void removeAutoGauge(final String pipeID) { Metric.PIPE_TSFILE_EVENT_TRANSFER_TIME.toString(), Tag.NAME.toString(), operator.getPipeName()); - remainingEventAndTimeOperatorMap.remove(pipeID); } //////////////////////////// register & deregister (pipe integration) //////////////////////////// @@ -202,18 +229,93 @@ private void removeAutoGauge(final String pipeID) { public void register(final IoTDBDataRegionSource source) { // The metric is global thus the regionId is omitted final String pipeID = source.getPipeName() + "_" + source.getCreationTime(); - remainingEventAndTimeOperatorMap - .computeIfAbsent( + final PipeDataNodeRemainingEventAndTimeOperator remainingOperator = + remainingEventAndTimeOperatorMap.computeIfAbsent( pipeID, k -> new PipeDataNodeRemainingEventAndTimeOperator( - source.getPipeName(), source.getCreationTime())) - .register(source); + source.getPipeName(), source.getCreationTime())); + remainingOperator.register(source); + final PipeDataNodeCompletionOperator completionOperator = + completionOperatorMap.computeIfAbsent( + pipeID, + k -> + new PipeDataNodeCompletionOperator( + () -> remainingOperator.getRemainingNonHeartbeatEvents() == 0, + () -> + PipeDataNodeAgent.task() + .getPipeCompletionSnapshot( + source.getPipeName(), source.getCreationTime()))); + if (source.getRealtimeSourceForCompletion() != null) { + completionOperator.registerDataRegionSource(source.getRealtimeSourceForCompletion()); + } if (Objects.nonNull(metricService)) { createMetrics(pipeID); } } + public void register( + final PipeRealtimeDataRegionSource source, final PipeDataRegionAssigner assigner) { + final String pipeID = source.getPipeName() + "_" + source.getCreationTime(); + final PipeDataNodeCompletionOperator operator = completionOperatorMap.get(pipeID); + if (operator != null) { + operator.register(source, assigner); + } + } + + public void deregister(final IoTDBDataRegionSource source) { + final String pipeID = source.getPipeName() + "_" + source.getCreationTime(); + final PipeDataNodeRemainingEventAndTimeOperator remainingOperator = + remainingEventAndTimeOperatorMap.get(pipeID); + if (remainingOperator != null) { + remainingOperator.deregister(source); + } + + final PipeDataNodeCompletionOperator operator = completionOperatorMap.get(pipeID); + if (operator != null && source.getRealtimeSourceForCompletion() != null) { + operator.deregisterDataRegionSource(source.getRealtimeSourceForCompletion()); + } + } + + public void deregister( + final PipeRealtimeDataRegionSource source, final PipeDataRegionAssigner assigner) { + final PipeDataNodeCompletionOperator operator = + completionOperatorMap.get(source.getPipeName() + "_" + source.getCreationTime()); + if (operator != null) { + operator.deregister(source, assigner); + } + } + + public void markDataRegionCompleted( + final String pipeName, + final long creationTime, + final int dataRegionId, + final PipeTaskMeta pipeTaskMeta, + final long assignerEpoch, + final long generation, + final long completionSourceId, + final CommitterKey committerKey) { + final PipeDataNodeCompletionOperator operator = + completionOperatorMap.get(pipeName + "_" + creationTime); + if (operator != null) { + operator.markCompleted( + dataRegionId, pipeTaskMeta, assignerEpoch, generation, completionSourceId, committerKey); + } + } + + public void markDataRegionInvalid( + final String pipeName, + final long creationTime, + final int dataRegionId, + final PipeTaskMeta pipeTaskMeta, + final long completionSourceId) { + final PipeDataNodeCompletionOperator operator = + completionOperatorMap.get(pipeName + "_" + creationTime); + if (operator != null) { + operator.markInvalid(dataRegionId, pipeTaskMeta, completionSourceId); + } + } + public void register(final IoTDBSchemaRegionSource source) { // The metric is global thus the regionId is omitted final String pipeID = source.getPipeName() + "_" + source.getCreationTime(); @@ -322,11 +424,11 @@ public void increaseHeartbeatEventCount(final String pipeName, final long creati } public void decreaseHeartbeatEventCount(final String pipeName, final long creationTime) { - remainingEventAndTimeOperatorMap - .computeIfAbsent( - pipeName + "_" + creationTime, - k -> new PipeDataNodeRemainingEventAndTimeOperator(pipeName, creationTime)) - .decreaseHeartbeatEventCount(); + final PipeDataNodeRemainingEventAndTimeOperator operator = + remainingEventAndTimeOperatorMap.get(pipeName + "_" + creationTime); + if (operator != null) { + operator.decreaseHeartbeatEventCount(); + } } public void thawRate(final String pipeID) { @@ -352,11 +454,11 @@ public void freezeRate(final String pipeID) { public void deregister(final String pipeID) { if (!remainingEventAndTimeOperatorMap.containsKey(pipeID)) { LOGGER.warn(DataNodePipeMessages.FAILED_TO_DEREGISTER_PIPE_REMAINING_EVENT_AND, pipeID); - return; - } - if (Objects.nonNull(metricService)) { + } else if (Objects.nonNull(metricService)) { removeMetrics(pipeID); } + remainingEventAndTimeOperatorMap.remove(pipeID); + completionOperatorMap.remove(pipeID); } public void markRegionCommit(final String pipeID, final boolean isDataRegion) { diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/IoTDBDataRegionSource.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/IoTDBDataRegionSource.java index 3e9c8ea4f7e9d..5e41c260d8fbe 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/IoTDBDataRegionSource.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/IoTDBDataRegionSource.java @@ -63,6 +63,7 @@ import java.time.ZoneId; import java.util.Arrays; import java.util.Objects; +import java.util.concurrent.atomic.AtomicBoolean; import java.util.concurrent.atomic.AtomicReference; import static org.apache.iotdb.commons.pipe.config.constant.PipeSourceConstant.EXTRACTOR_DATABASE_KEY; @@ -137,6 +138,7 @@ public class IoTDBDataRegionSource extends IoTDBSource { private boolean hasNoExtractionNeed = true; private boolean shouldExtractDeletion = false; + private final AtomicBoolean completionSourceStarted = new AtomicBoolean(false); @Override public void validate(final PipeParameterValidator validator) throws Exception { @@ -549,6 +551,7 @@ public void start() throws Exception { dataRegionIdObject, () -> startHistoricalExtractorAndRealtimeExtractor(exceptionHolder))) { rethrowExceptionIfAny(exceptionHolder); + completionSourceStarted.set(true); LOGGER.info( DataNodePipeMessages.PIPE_STARTED_HISTORICAL_SOURCE_AND_REALTIME_SOURCE, @@ -624,6 +627,8 @@ public Event supply() throws Exception { @Override public void close() throws Exception { + completionSourceStarted.set(false); + PipeDataNodeSinglePipeMetrics.getInstance().deregister(this); if (hasNoExtractionNeed || !hasBeenStarted.get()) { return; } @@ -635,6 +640,16 @@ public void close() throws Exception { } } + public PipeRealtimeDataRegionSource getRealtimeSourceForCompletion() { + return realtimeSource; + } + + public boolean isReadyForCompletion() { + return completionSourceStarted.get() + && Objects.nonNull(historicalSource) + && historicalSource.hasConsumedAll(); + } + //////////////////////////// APIs provided for metric framework //////////////////////////// public int getHistoricalTsFileInsertionEventCount() { 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..0c4c9b1d9d477 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 @@ -207,13 +207,17 @@ private void markTsFileEpochActive(final TsFileEpoch tsFileEpoch) { private void markTsFileEpochDegraded(final TsFileEpoch tsFileEpoch) { activeTsFileEpochs.add(tsFileEpoch); - degradedTsFileEpochs.add(tsFileEpoch); + if (degradedTsFileEpochs.add(tsFileEpoch)) { + markCompletionStateChanged(); + } reportTsFileEpochDegradedStatus(); } private void clearTsFileEpoch(final TsFileEpoch tsFileEpoch) { activeTsFileEpochs.remove(tsFileEpoch); - degradedTsFileEpochs.remove(tsFileEpoch); + if (degradedTsFileEpochs.remove(tsFileEpoch)) { + markCompletionStateChanged(); + } reportTsFileEpochDegradedStatus(); } @@ -227,12 +231,20 @@ private void reportTsFileEpochDegradedStatus() { } } + @Override + public boolean isTsFileEpochDegraded() { + return !degradedTsFileEpochs.isEmpty(); + } + @Override public void close() throws Exception { try { super.close(); } finally { activeTsFileEpochs.clear(); + if (!degradedTsFileEpochs.isEmpty()) { + markCompletionStateChanged(); + } degradedTsFileEpochs.clear(); PipeDataNodeAgent.task().clearPipeTsFileEpochDegraded(pipeName, creationTime, dataRegionId); } @@ -346,6 +358,7 @@ private Event supplyTsFileInsertion(final PipeRealtimeEvent event) { DataNodePipeMessages.EVENT_CAN_NOT_BE_SUPPLIED_BECAUSE_DATA_IS_LOST, event.getEvent()); LOGGER.error(errorMessage); + markDataRegionCompletionInvalid(); PipeDataNodeAgent.runtime() .report(pipeTaskMeta, new PipeRuntimeNonCriticalException(errorMessage)); PipeTsFileEpochProgressIndexKeeper.getInstance() diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionLogSource.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionLogSource.java index a90441bae2b10..21325d8222414 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionLogSource.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionLogSource.java @@ -117,6 +117,7 @@ public Event supply() { DataNodePipeMessages.EVENT_CAN_NOT_BE_SUPPLIED_BECAUSE_DATA_IS_LOST, realtimeEvent.getEvent()); LOGGER.error(errorMessage); + markDataRegionCompletionInvalid(); PipeDataNodeAgent.runtime() .report(pipeTaskMeta, new PipeRuntimeNonCriticalException(errorMessage)); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionSource.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionSource.java index 732275824eb46..ee0686c08ca83 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionSource.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionSource.java @@ -42,6 +42,7 @@ import org.apache.iotdb.db.pipe.event.common.tablet.PipeInsertNodeTabletInsertionEvent; 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.metric.overview.PipeDataNodeSinglePipeMetrics; import org.apache.iotdb.db.pipe.metric.source.PipeDataRegionEventCounter; import org.apache.iotdb.db.pipe.processor.iotconsensusv2.IoTConsensusV2Processor; import org.apache.iotdb.db.pipe.source.dataregion.DataRegionListeningFilter; @@ -94,11 +95,14 @@ public abstract class PipeRealtimeDataRegionSource implements PipeExtractor { private static final Logger LOGGER = LoggerFactory.getLogger(PipeRealtimeDataRegionSource.class); + private static final AtomicLong COMPLETION_SOURCE_ID_GENERATOR = new AtomicLong(0); protected String pipeName; protected long creationTime; protected int dataRegionId = -1; protected PipeTaskMeta pipeTaskMeta; + private final long completionSourceId = COMPLETION_SOURCE_ID_GENERATOR.incrementAndGet(); + private final AtomicLong completionStateVersion = new AtomicLong(0); protected boolean shouldExtractInsertion; protected boolean shouldExtractDeletion; @@ -414,13 +418,15 @@ public final void extract(final PipeRealtimeEvent event) { protected void extractHeartbeat(final PipeRealtimeEvent event) { // Record the pending queue size before trying to put heartbeatEvent into queue - ((PipeHeartbeatEvent) event.getEvent()).recordExtractorQueueSize(pendingQueue); + final PipeHeartbeatEvent heartbeatEvent = (PipeHeartbeatEvent) event.getEvent(); + heartbeatEvent.recordExtractorQueueSize(pendingQueue); final Event lastEvent = pendingQueue.peekLast(); - if (lastEvent instanceof PipeRealtimeEvent + if (!heartbeatEvent.isCompletionBarrier() + && lastEvent instanceof PipeRealtimeEvent && ((PipeRealtimeEvent) lastEvent).getEvent() instanceof PipeHeartbeatEvent && (((PipeHeartbeatEvent) ((PipeRealtimeEvent) lastEvent).getEvent()).isShouldPrintMessage() - || !((PipeHeartbeatEvent) event.getEvent()).isShouldPrintMessage())) { + || !heartbeatEvent.isShouldPrintMessage())) { // If the last event in the pending queue is a heartbeat event, we should not extract any more // heartbeat events to avoid OOM when the pipe is stopped. // Besides, the printable event has higher priority to stay in queue to enable metrics report. @@ -438,7 +444,9 @@ protected void extractProgressReportEvent(final PipeRealtimeEvent event) { // yet while (true) { final PipeRealtimeEvent lastEvent = ((PipeRealtimeEvent) pendingQueue.peekLast()); - if (lastEvent == null || !(lastEvent.getEvent() instanceof PipeHeartbeatEvent)) { + if (lastEvent == null + || !(lastEvent.getEvent() instanceof PipeHeartbeatEvent) + || ((PipeHeartbeatEvent) lastEvent.getEvent()).isCompletionBarrier()) { break; } final PipeRealtimeEvent droppedEvent = (PipeRealtimeEvent) pendingQueue.pollLast(); @@ -525,6 +533,7 @@ protected Event supplyDirectly(final PipeRealtimeEvent event) { DataNodePipeMessages.EVENT_CAN_NOT_BE_SUPPLIED_BECAUSE_DATA_IS_LOST, event.getEvent()); LOGGER.error(errorMessage); + markDataRegionCompletionInvalid(); PipeDataNodeAgent.runtime() .report(pipeTaskMeta, new PipeRuntimeNonCriticalException(errorMessage)); return null; @@ -579,6 +588,29 @@ public final int getDataRegionId() { return dataRegionId; } + public final long getCompletionSourceId() { + return completionSourceId; + } + + public final long getCompletionStateVersion() { + return completionStateVersion.get(); + } + + protected final void markCompletionStateChanged() { + completionStateVersion.incrementAndGet(); + } + + protected final void markDataRegionCompletionInvalid() { + PipeDataNodeSinglePipeMetrics.getInstance() + .markDataRegionInvalid( + pipeName, creationTime, dataRegionId, pipeTaskMeta, completionSourceId); + } + + /** Whether this source is waiting for a TsFile to recover discarded tablet events. */ + public boolean isTsFileEpochDegraded() { + return false; + } + public final long getRealtimeDataExtractionStartTime() { return realtimeDataExtractionStartTime; } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionTsFileSource.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionTsFileSource.java index 73bef31d85bdf..077bd1897b6ab 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionTsFileSource.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionTsFileSource.java @@ -100,6 +100,7 @@ public Event supply() { DataNodePipeMessages.EVENT_CAN_NOT_BE_SUPPLIED_BECAUSE_DATA_IS_LOST, realtimeEvent.getEvent()); LOGGER.error(errorMessage); + markDataRegionCompletionInvalid(); PipeDataNodeAgent.runtime() .report(pipeTaskMeta, new PipeRuntimeNonCriticalException(errorMessage)); PipeTsFileEpochProgressIndexKeeper.getInstance() diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/assigner/PipeDataRegionAssigner.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/assigner/PipeDataRegionAssigner.java index 6e40be3d324bb..5840761871e41 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/assigner/PipeDataRegionAssigner.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/assigner/PipeDataRegionAssigner.java @@ -32,6 +32,7 @@ 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.metric.overview.PipeDataNodeSinglePipeMetrics; import org.apache.iotdb.db.pipe.metric.source.PipeAssignerMetrics; import org.apache.iotdb.db.pipe.metric.source.PipeDataRegionEventCounter; import org.apache.iotdb.db.pipe.source.dataregion.realtime.PipeRealtimeDataRegionSource; @@ -49,10 +50,14 @@ import java.io.Closeable; import java.util.Objects; import java.util.Set; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.locks.ReentrantLock; +import java.util.function.Supplier; public class PipeDataRegionAssigner implements Closeable { private static final Logger LOGGER = LoggerFactory.getLogger(PipeDataRegionAssigner.class); + private static final AtomicLong ASSIGNER_EPOCH_GENERATOR = new AtomicLong(0); /** * The {@link PipeDataRegionMatcher} is used to match the event with the source based on the @@ -64,6 +69,11 @@ public class PipeDataRegionAssigner implements Closeable { private final DisruptorQueue disruptor; private final int dataRegionId; + private final long assignerEpoch = ASSIGNER_EPOCH_GENERATOR.incrementAndGet(); + private final AtomicLong publishedDataGeneration = new AtomicLong(0); + private final AtomicLong publicationFailureEpoch = new AtomicLong(0); + private final AtomicLong completionInvalidationEpoch = new AtomicLong(0); + private final ReentrantLock publicationLock = new ReentrantLock(); private Boolean isTableModel; @@ -94,42 +104,189 @@ public PipeDataRegionAssigner(final int dataRegionId) { } public void publishToAssign(final PipeRealtimeEvent event) { + publicationLock.lock(); + try { + final EnrichedEvent innerEvent = event.getEvent(); + final boolean isDataEvent = !(innerEvent instanceof PipeHeartbeatEvent); + if (isDataEvent) { + publishedDataGeneration.incrementAndGet(); + } + publishToAssignInternal(event, isDataEvent); + } finally { + publicationLock.unlock(); + } + } + + public void publishDataEventToAssign(final Supplier eventSupplier) { + publishDataEventToAssign(eventSupplier, false); + } + + public void publishInsertDataEventToAssign(final Supplier eventSupplier) { + publishDataEventToAssign(eventSupplier, true); + } + + private void publishDataEventToAssign( + final Supplier eventSupplier, final boolean invalidateCompletionBarrier) { + publicationLock.lock(); + try { + publishedDataGeneration.incrementAndGet(); + if (invalidateCompletionBarrier) { + completionInvalidationEpoch.incrementAndGet(); + } + final PipeRealtimeEvent event; + try { + event = eventSupplier.get(); + } catch (final RuntimeException | Error e) { + markPublicationFailed(); + throw e; + } + if (event != null) { + publishToAssignInternal(event, true); + } else { + markPublicationFailed(); + } + } finally { + publicationLock.unlock(); + } + } + + private void publishToAssignInternal(final PipeRealtimeEvent event, final boolean isDataEvent) { + final EnrichedEvent innerEvent = event.getEvent(); if (!event.increaseReferenceCount(PipeDataRegionAssigner.class.getName())) { + if (isDataEvent) { + markPublicationFailed(); + } LOGGER.warn(DataNodePipeMessages.THE_REFERENCE_COUNT_OF_THE_REALTIME_EVENT, event); return; } - final EnrichedEvent innerEvent = event.getEvent(); eventCounter.increaseEventCount(innerEvent); if (innerEvent instanceof PipeHeartbeatEvent) { ((PipeHeartbeatEvent) innerEvent).onPublished(); } - synchronized (this) { - if (disruptor.isClosed()) { - onAssignedHook(event); - return; - } - inFlightPublishCount++; - } - + boolean shouldReleaseDirectly = false; boolean isPublished = false; try { - isPublished = disruptor.publishOrDrop(event); - } finally { + if (innerEvent instanceof PipeHeartbeatEvent) { + final PipeHeartbeatEvent heartbeatEvent = (PipeHeartbeatEvent) innerEvent; + if (heartbeatEvent.isCompletionBarrier()) { + heartbeatEvent.bindCompletionBarrier(assignerEpoch, publishedDataGeneration.get()); + } + } + synchronized (this) { - inFlightPublishCount--; - if (inFlightPublishCount == 0) { - notifyAll(); + if (disruptor.isClosed()) { + shouldReleaseDirectly = true; + } else { + inFlightPublishCount++; + } + } + + if (!shouldReleaseDirectly) { + try { + isPublished = disruptor.publishOrDrop(event); + } finally { + synchronized (this) { + inFlightPublishCount--; + if (inFlightPublishCount == 0) { + notifyAll(); + } + } } } + } catch (final RuntimeException | Error e) { + if (isDataEvent) { + markPublicationFailed(); + } + throw e; } - if (!isPublished) { + if (shouldReleaseDirectly || !isPublished) { + if (isDataEvent) { + markPublicationFailed(); + } onAssignedHook(event); } } + private void markPublicationFailed() { + publicationFailureEpoch.incrementAndGet(); + } + + /** + * Advances the published data generation for a data event that was ignored before publication. + * This deliberately keeps the current completion token valid: full-flush close callbacks may be + * ignored when no source listens to TsFile events, but their generations still need to be covered + * by the barrier for that flush. + */ + public void invalidateCompletion() { + invalidateCompletion(false); + } + + /** + * Advances the published data generation and invalidates the current completion token for an + * ignored insert. The insert may create a working TsFile processor that the in-progress full + * flush does not cover, so its barrier must not be published. + */ + public void invalidateCompletionBarrier() { + invalidateCompletion(true); + } + + private void invalidateCompletion(final boolean invalidateCompletionBarrier) { + publicationLock.lock(); + try { + publishedDataGeneration.incrementAndGet(); + if (invalidateCompletionBarrier) { + completionInvalidationEpoch.incrementAndGet(); + } + } finally { + publicationLock.unlock(); + } + } + + /** + * Atomically advances the generation, invalidates any prior token, and returns the token for a + * new full flush. The corresponding completion barrier is accepted only if no insert invalidates + * this token before the flush finishes. + */ + public CompletionToken invalidateCompletionAndGetToken() { + publicationLock.lock(); + try { + publishedDataGeneration.incrementAndGet(); + return new CompletionToken(assignerEpoch, completionInvalidationEpoch.incrementAndGet()); + } finally { + publicationLock.unlock(); + } + } + + public boolean publishCompletionBarrier(final CompletionToken token) { + publicationLock.lock(); + try { + if (token == null + || token.assignerEpoch != assignerEpoch + || token.completionInvalidationEpoch != completionInvalidationEpoch.get()) { + return false; + } + publishToAssignInternal( + PipeRealtimeEventFactory.createCompletionBarrierEvent(dataRegionId), false); + return true; + } finally { + publicationLock.unlock(); + } + } + + public static final class CompletionToken { + + private final long assignerEpoch; + private final long completionInvalidationEpoch; + + private CompletionToken(final long assignerEpoch, final long completionInvalidationEpoch) { + this.assignerEpoch = assignerEpoch; + this.completionInvalidationEpoch = completionInvalidationEpoch; + } + } + private void onAssignedHook(final PipeRealtimeEvent realtimeEvent) { realtimeEvent.gcSchemaInfo(); realtimeEvent.decreaseReferenceCount(PipeDataRegionAssigner.class.getName(), false); @@ -148,6 +305,17 @@ private void assignToSource( return; } + try { + assignToSourceInternal(event); + } catch (final RuntimeException | Error e) { + if (!(event.getEvent() instanceof PipeHeartbeatEvent)) { + markPublicationFailed(); + } + throw e; + } + } + + private void assignToSourceInternal(final PipeRealtimeEvent event) { final Pair, Set> matchedAndUnmatched = matcher.match(event); @@ -165,6 +333,7 @@ private void assignToSource( source.getPipeName(), source.getCreationTime(), source.getPipeTaskMeta()); reportEvent.bindProgressIndex(event.getProgressIndex()); if (!reportEvent.increaseReferenceCount(PipeDataRegionAssigner.class.getName())) { + markPublicationFailed(); LOGGER.warn( DataNodePipeMessages.THE_REFERENCE_COUNT_OF_THE_EVENT_CANNOT, reportEvent); return; @@ -188,6 +357,12 @@ private void assignToSource( source.getRealtimeDataExtractionEndTime()); final EnrichedEvent innerEvent = copiedEvent.getEvent(); + if (innerEvent instanceof PipeHeartbeatEvent + && ((PipeHeartbeatEvent) innerEvent).isCompletionBarrier()) { + ((PipeHeartbeatEvent) innerEvent) + .bindCompletionSource(source.getCompletionSourceId()); + } + if (innerEvent instanceof PipeTsFileInsertionEvent) { final PipeTsFileInsertionEvent tsFileInsertionEvent = (PipeTsFileInsertionEvent) innerEvent; @@ -211,6 +386,9 @@ private void assignToSource( } if (!copiedEvent.increaseReferenceCount(PipeDataRegionAssigner.class.getName())) { + if (!(event.getEvent() instanceof PipeHeartbeatEvent)) { + markPublicationFailed(); + } LOGGER.warn( DataNodePipeMessages.THE_REFERENCE_COUNT_OF_THE_EVENT_CANNOT, copiedEvent); return; @@ -234,6 +412,7 @@ private void assignToSource( source.getPipeName(), source.getCreationTime(), source.getPipeTaskMeta()); reportEvent.bindProgressIndex(event.getProgressIndex()); if (!reportEvent.increaseReferenceCount(PipeDataRegionAssigner.class.getName())) { + markPublicationFailed(); LOGGER.warn( DataNodePipeMessages.THE_REFERENCE_COUNT_OF_THE_EVENT_CANNOT, reportEvent); return; @@ -244,7 +423,13 @@ private void assignToSource( } public synchronized void startAssignTo(final PipeRealtimeDataRegionSource source) { - matcher.register(source); + PipeDataNodeSinglePipeMetrics.getInstance().register(source, this); + try { + matcher.register(source); + } catch (final RuntimeException | Error e) { + PipeDataNodeSinglePipeMetrics.getInstance().deregister(source, this); + throw e; + } if (source.isNeedListenToTsFile()) { listenToTsFileSourceCount++; } @@ -256,6 +441,7 @@ public synchronized void startAssignTo(final PipeRealtimeDataRegionSource source public synchronized void stopAssignTo(final PipeRealtimeDataRegionSource source) { matcher.deregister(source); + PipeDataNodeSinglePipeMetrics.getInstance().deregister(source, this); if (source.isNeedListenToTsFile()) { listenToTsFileSourceCount--; } @@ -348,4 +534,16 @@ private void logSourceAssignmentChange( public Boolean isTableModel() { return isTableModel; } + + public long getAssignerEpoch() { + return assignerEpoch; + } + + public long getPublishedDataGeneration() { + return publishedDataGeneration.get(); + } + + public long getPublicationFailureEpoch() { + return publicationFailureEpoch.get(); + } } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/listener/PipeInsertionDataNodeListener.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/listener/PipeInsertionDataNodeListener.java index 7e3c25e806293..9c2a4bee090c5 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/listener/PipeInsertionDataNodeListener.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/listener/PipeInsertionDataNodeListener.java @@ -26,6 +26,7 @@ import org.apache.iotdb.db.pipe.event.realtime.PipeRealtimeEventFactory; import org.apache.iotdb.db.pipe.source.dataregion.realtime.PipeRealtimeDataRegionSource; import org.apache.iotdb.db.pipe.source.dataregion.realtime.assigner.PipeDataRegionAssigner; +import org.apache.iotdb.db.pipe.source.dataregion.realtime.assigner.PipeDataRegionAssigner.CompletionToken; import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.AbstractDeleteDataNode; import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.InsertNode; import org.apache.iotdb.db.queryengine.plan.planner.plan.node.write.RelationalDeleteDataNode; @@ -34,6 +35,7 @@ import java.util.Objects; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.ConcurrentMap; +import java.util.concurrent.atomic.AtomicReference; /** * PipeInsertionEventListener is a singleton in each data node. @@ -101,14 +103,19 @@ public void listenToTsFile( final boolean isLoaded) { final PipeDataRegionAssigner assigner = dataRegionId2Assigner.get(dataRegionId); + if (assigner == null) { + return; + } // only events from registered data region with tsfile listeners will be extracted - if (assigner == null || !assigner.shouldListenToTsFile()) { + if (!assigner.shouldListenToTsFile()) { + assigner.invalidateCompletion(); return; } - assigner.publishToAssign( - PipeRealtimeEventFactory.createRealtimeEvent( - assigner.isTableModel(), databaseName, tsFileResource, isLoaded)); + assigner.publishDataEventToAssign( + () -> + PipeRealtimeEventFactory.createRealtimeEvent( + assigner.isTableModel(), databaseName, tsFileResource, isLoaded)); } public void listenToInsertNode( @@ -118,14 +125,19 @@ public void listenToInsertNode( final TsFileResource tsFileResource) { final PipeDataRegionAssigner assigner = dataRegionId2Assigner.get(dataRegionId); + if (assigner == null) { + return; + } // only events from registered data region with insert listeners will be extracted - if (assigner == null || !assigner.shouldListenToInsertNode()) { + if (!assigner.shouldListenToInsertNode()) { + assigner.invalidateCompletionBarrier(); return; } - assigner.publishToAssign( - PipeRealtimeEventFactory.createRealtimeEvent( - assigner.isTableModel(), databaseName, insertNode, tsFileResource)); + assigner.publishInsertDataEventToAssign( + () -> + PipeRealtimeEventFactory.createRealtimeEvent( + assigner.isTableModel(), databaseName, insertNode, tsFileResource)); } public DeletionResource listenToDeleteData( @@ -137,25 +149,25 @@ public DeletionResource listenToDeleteData( && ((RelationalDeleteDataNode) node).getModEntries().isEmpty()) { return null; } + final AtomicReference deletionResourceHolder = new AtomicReference<>(); + assigner.publishDataEventToAssign( + () -> { + // Register a deletionResource and return it to DataRegion. A deleteNode generated by a + // remote consensus leader should not be persisted to DAL. + final DeletionResourceManager manager = DeletionResourceManager.getInstance(regionId); + if (Objects.nonNull(manager) + && DeletionResource.isDeleteNodeGeneratedInLocalByIoTV2(node)) { + final DeletionResource deletionResource = manager.registerDeletionResource(node); + deletionResourceHolder.set(deletionResource); + // If persistence failed, skip this event to stay consistent with storage engine. + if (deletionResource.waitForResult() == DeletionResource.Status.FAILURE) { + return null; + } + } + return PipeRealtimeEventFactory.createRealtimeEvent(node); + }); - final DeletionResource deletionResource; - // register a deletionResource and return it to DataRegion - final DeletionResourceManager manager = DeletionResourceManager.getInstance(regionId); - // deleteNode generated by remote consensus leader shouldn't be persisted to DAL. - if (Objects.nonNull(manager) && DeletionResource.isDeleteNodeGeneratedInLocalByIoTV2(node)) { - deletionResource = manager.registerDeletionResource(node); - // if persist failed, skip sending/publishing this event to keep consistency with the - // behavior of storage engine. - if (deletionResource.waitForResult() == DeletionResource.Status.FAILURE) { - return deletionResource; - } - } else { - deletionResource = null; - } - - assigner.publishToAssign(PipeRealtimeEventFactory.createRealtimeEvent(node)); - - return deletionResource; + return deletionResourceHolder.get(); } public void listenToHeartbeat(final boolean shouldPrintMessage) { @@ -165,6 +177,19 @@ public void listenToHeartbeat(final boolean shouldPrintMessage) { PipeRealtimeEventFactory.createRealtimeEvent(key, shouldPrintMessage))); } + public CompletionToken invalidateDataRegionCompletion(final int dataRegionId) { + final PipeDataRegionAssigner assigner = dataRegionId2Assigner.get(dataRegionId); + return assigner == null ? null : assigner.invalidateCompletionAndGetToken(); + } + + public void listenToCompletionBarrier( + final int dataRegionId, final CompletionToken completionToken) { + final PipeDataRegionAssigner assigner = dataRegionId2Assigner.get(dataRegionId); + if (assigner != null) { + assigner.publishCompletionBarrier(completionToken); + } + } + public boolean isEmpty() { return dataRegionId2Assigner.isEmpty(); } diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java index 2062ba709784c..8a10ec4379de9 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/StorageEngine.java @@ -58,6 +58,8 @@ import org.apache.iotdb.db.exception.runtime.StorageEngineFailureException; import org.apache.iotdb.db.i18n.StorageEngineMessages; import org.apache.iotdb.db.pipe.agent.PipeDataNodeAgent; +import org.apache.iotdb.db.pipe.source.dataregion.realtime.assigner.PipeDataRegionAssigner.CompletionToken; +import org.apache.iotdb.db.pipe.source.dataregion.realtime.listener.PipeInsertionDataNodeListener; import org.apache.iotdb.db.queryengine.plan.analyze.cache.schema.DataNodeTTLCache; import org.apache.iotdb.db.queryengine.plan.planner.plan.node.load.LoadTsFilePieceNode; import org.apache.iotdb.db.queryengine.plan.scheduler.load.LoadTsFileScheduler; @@ -522,7 +524,7 @@ public void syncCloseAllProcessor() { tasks.add( cachedThreadPool.submit( () -> { - dataRegion.syncCloseAllWorkingTsFileProcessors(); + syncCloseAllProcessorAndPublishCompletionBarrier(dataRegion); return null; })); } @@ -553,7 +555,7 @@ public void syncCloseProcessorsInDatabase(String databaseName) { tasks.add( cachedThreadPool.submit( () -> { - dataRegion.syncCloseAllWorkingTsFileProcessors(); + syncCloseAllProcessorAndPublishCompletionBarrier(dataRegion); return null; })); } @@ -568,7 +570,7 @@ public void syncCloseProcessorsInRegion(List dataRegionIds) { tasks.add( cachedThreadPool.submit( () -> { - dataRegion.syncCloseAllWorkingTsFileProcessors(); + syncCloseAllProcessorAndPublishCompletionBarrier(dataRegion); return null; })); } @@ -609,6 +611,15 @@ private void checkResults(List> tasks, String errorMsg) { } } + private void syncCloseAllProcessorAndPublishCompletionBarrier(final DataRegion dataRegion) + throws InterruptedException, ExecutionException { + final PipeInsertionDataNodeListener listener = PipeInsertionDataNodeListener.getInstance(); + final int dataRegionId = dataRegion.getDataRegionId(); + final CompletionToken completionToken = listener.invalidateDataRegionCompletion(dataRegionId); + dataRegion.syncCloseAllWorkingAndClosingTsFileProcessors(); + listener.listenToCompletionBarrier(dataRegionId, completionToken); + } + /** * merge all databases. * diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java index d2619a9de65bf..9ad92dec35a45 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/DataRegion.java @@ -2538,6 +2538,78 @@ public void syncCloseAllWorkingTsFileProcessors() { } } + /** + * Closes all working TsFile processors and waits for processors that were already closing. + * + *

This stronger variant is used by a full flush before a Pipe completion barrier is published. + * Unlike {@link #syncCloseAllWorkingTsFileProcessors()}, failures are propagated to the caller so + * a failed flush can never publish a completion barrier. + */ + public void syncCloseAllWorkingAndClosingTsFileProcessors() + throws InterruptedException, ExecutionException { + final Set targets = new HashSet<>(); + final Set targetsToClose = new HashSet<>(); + final List> closingFutures = new ArrayList<>(); + + writeLock("syncCloseAllWorkingAndClosingTsFileProcessors"); + try { + targets.addAll(closingSequenceTsFileProcessor); + targets.addAll(closingUnSequenceTsFileProcessor); + for (final TsFileProcessor target : targets) { + closingFutures.add(target.getCloseFuture()); + } + + targetsToClose.addAll(workSequenceTsFileProcessors.values()); + targetsToClose.addAll(workUnsequenceTsFileProcessors.values()); + targets.addAll(targetsToClose); + } finally { + writeUnlock(); + } + + int closedProcessorCount = 0; + while (!targetsToClose.isEmpty()) { + final List> currentFutures = new ArrayList<>(); + final Set targetsWaitingForOrdinaryFlush = new HashSet<>(); + + writeLock("syncCloseAllWorkingAndClosingTsFileProcessors"); + try { + final Iterator iterator = targetsToClose.iterator(); + while (iterator.hasNext()) { + final TsFileProcessor target = iterator.next(); + final Future future = + asyncCloseOneTsFileProcessorForFullFlush(target.isSequence(), target); + currentFutures.add(future); + if (target.alreadyMarkedClosing()) { + closingFutures.add(future); + iterator.remove(); + closedProcessorCount++; + } else { + targetsWaitingForOrdinaryFlush.add(target); + } + } + } finally { + writeUnlock(); + } + + for (final Future future : currentFutures) { + if (future != null) { + future.get(); + } + } + for (final TsFileProcessor target : targetsWaitingForOrdinaryFlush) { + target.waitUntilFlushManagerReleased(); + } + } + + WritingMetrics.getInstance().recordManualFlushMemTableCount(closedProcessorCount); + for (final Future future : closingFutures) { + if (future != null) { + future.get(); + } + } + waitClosingTsFileProcessorFinished(targets); + } + public void syncCloseWorkingTsFileProcessors(boolean sequence) { try { writeLock("syncCloseWorkingTsFileProcessors"); @@ -2602,6 +2674,27 @@ private void waitClosingTsFileProcessorFinished() throws InterruptedException { } } + private void waitClosingTsFileProcessorFinished(final Set targets) + throws InterruptedException { + while (containsClosingTarget(targets)) { + synchronized (closeStorageGroupCondition) { + if (containsClosingTarget(targets)) { + closeStorageGroupCondition.wait(60_000); + } + } + } + } + + private boolean containsClosingTarget(final Set targets) { + for (final TsFileProcessor target : targets) { + if (closingSequenceTsFileProcessor.contains(target) + || closingUnSequenceTsFileProcessor.contains(target)) { + return true; + } + } + return false; + } + /** close all working tsfile processors */ public List> asyncCloseAllWorkingTsFileProcessors() { writeLock("asyncCloseAllWorkingTsFileProcessors"); @@ -2630,6 +2723,53 @@ public List> asyncCloseAllWorkingTsFileProcessors() { return futures; } + private Future asyncCloseOneTsFileProcessorForFullFlush( + final boolean sequence, final TsFileProcessor tsFileProcessor) { + if (tsFileProcessor == null) { + return CompletableFuture.completedFuture(null); + } + if (tsFileProcessor.alreadyMarkedClosing()) { + return tsFileProcessor.getCloseFuture(); + } + + final Set closingTsFileProcessors = + sequence ? closingSequenceTsFileProcessor : closingUnSequenceTsFileProcessor; + final TreeMap workTsFileProcessors = + sequence ? workSequenceTsFileProcessors : workUnsequenceTsFileProcessors; + + // Register before scheduling so even a very fast close callback can remove this processor. + closingTsFileProcessors.add(tsFileProcessor); + final Future future = tsFileProcessor.asyncCloseForFullFlush(); + if (!tsFileProcessor.alreadyMarkedClosing()) { + // An ordinary flush is still running. Keep the processor working, wait outside the locks, + // and retry after the flush manager has released it. + closingTsFileProcessors.remove(tsFileProcessor); + return future; + } + + if (future != null && future.isDone()) { + closingTsFileProcessors.remove(tsFileProcessor); + } + final boolean removedWorkingProcessor = + workTsFileProcessors.get(tsFileProcessor.getTimeRangeId()) == tsFileProcessor; + if (removedWorkingProcessor) { + workTsFileProcessors.remove(tsFileProcessor.getTimeRangeId()); + } + + final TsFileResource resource = tsFileProcessor.getTsFileResource(); + logger.info( + StorageEngineMessages.STORAGE_LOG_ASYNC_CLOSE_TSFILE_FILE_START_TIME_FILE_END_TIME_65020832, + resource.getTsFile().getAbsolutePath(), + resource.getFileStartTime(), + resource.getFileEndTime()); + if (removedWorkingProcessor + && workSequenceTsFileProcessors.get(tsFileProcessor.getTimeRangeId()) == null + && workUnsequenceTsFileProcessors.get(tsFileProcessor.getTimeRangeId()) == null) { + WritingMetrics.getInstance().recordActiveTimePartitionCount(-1); + } + return future; + } + /** force close all working tsfile processors */ public void forceCloseAllWorkingTsFileProcessors() throws TsFileProcessorException { writeLock("forceCloseAllWorkingTsFileProcessors"); diff --git a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java index 4528b44f59894..6c03ca79dcb3a 100644 --- a/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java +++ b/iotdb-core/datanode/src/main/java/org/apache/iotdb/db/storageengine/dataregion/memtable/TsFileProcessor.java @@ -1344,10 +1344,28 @@ public void syncClose() throws ExecutionException { /** async close one tsfile, register and close it by another thread */ public Future asyncClose() { + return asyncClose(false); + } + + /** + * Tries to close this TsFile after any previous ordinary flush has finished. + * + *

If an ordinary flush is still being managed by the flush manager, this method only returns + * that flush's future. The caller must wait for it and retry until {@link + * #alreadyMarkedClosing()} becomes {@code true}. This avoids adding a close signal while the old + * flush task is about to leave the flush manager. + */ + public Future asyncCloseForFullFlush() { + return asyncClose(true); + } + + private Future asyncClose(final boolean ignorePreviousFlushFuture) { flushQueryLock.writeLock().lock(); logFlushQueryWriteLocked(); try { - if (closeFuture != null) { + if (shouldClose + || managedByFlushManager + || (!ignorePreviousFlushFuture && closeFuture != null)) { return closeFuture; } @@ -1398,12 +1416,12 @@ public Future asyncClose() { dataRegionName, tsFileResource.getTsFile().getName(), e); + return CompletableFuture.failedFuture(e); } } finally { flushQueryLock.writeLock().unlock(); logFlushQueryWriteUnlocked(); } - return CompletableFuture.completedFuture(null); } /** Put the working memtable into flushing list and set the working memtable to null */ @@ -1844,9 +1862,27 @@ public boolean isManagedByFlushManager() { } public void setManagedByFlushManager(boolean managedByFlushManager) { - this.managedByFlushManager = managedByFlushManager; + flushQueryLock.writeLock().lock(); + try { + this.managedByFlushManager = managedByFlushManager; + if (!managedByFlushManager) { + closeFuture = CompletableFuture.completedFuture(null); + } + } finally { + flushQueryLock.writeLock().unlock(); + } if (!managedByFlushManager) { - closeFuture = CompletableFuture.completedFuture(null); + synchronized (this) { + notifyAll(); + } + } + } + + public void waitUntilFlushManagerReleased() throws InterruptedException { + synchronized (this) { + while (managedByFlushManager) { + wait(); + } } } diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollectorTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollectorTest.java index 029a722c8a95c..de89e6c02ab47 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollectorTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/agent/task/connection/PipeEventCollectorTest.java @@ -53,7 +53,7 @@ private void verifyCollectorDoesNotOfferEventsOfDroppedPipe( pendingQueue.discardEventsOfPipe("pipe", 1L, 1); final PipeEventCollector droppedPipeCollector = - new PipeEventCollector(pendingQueue, 1L, 1, false, false, false); + new PipeEventCollector(pendingQueue, 1L, 1, Long.MIN_VALUE, false, false, false); final PipeRawTabletInsertionEvent droppedPipeEvent = createPipeRawTabletInsertionEvent("pipe", 1L); droppedPipeCollector.collect(droppedPipeEvent); @@ -62,7 +62,7 @@ private void verifyCollectorDoesNotOfferEventsOfDroppedPipe( Assert.assertEquals(0, pendingQueue.size()); final PipeEventCollector recreatedPipeCollector = - new PipeEventCollector(pendingQueue, 2L, 1, false, false, false); + new PipeEventCollector(pendingQueue, 2L, 1, Long.MIN_VALUE, false, false, false); final PipeRawTabletInsertionEvent recreatedPipeEvent = createPipeRawTabletInsertionEvent("pipe", 2L); recreatedPipeCollector.collect(recreatedPipeEvent); diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/common/heartbeat/PipeHeartbeatEventTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/common/heartbeat/PipeHeartbeatEventTest.java new file mode 100644 index 0000000000000..93e5b4702ea4f --- /dev/null +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/event/common/heartbeat/PipeHeartbeatEventTest.java @@ -0,0 +1,83 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.iotdb.db.pipe.event.common.heartbeat; + +import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta; +import org.apache.iotdb.commons.pipe.event.EnrichedEvent; + +import org.junit.Test; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; +import static org.mockito.Mockito.mock; + +public class PipeHeartbeatEventTest { + + @Test + public void testOnlyBoundCompletionBarrierNeedsOrderedCommit() { + final PipeHeartbeatEvent rawBarrier = new PipeHeartbeatEvent(1, false, true); + rawBarrier.bindCompletionBarrier(2L, 3L); + assertFalse(rawBarrier.needToCommit()); + + final PipeTaskMeta taskMeta = mock(PipeTaskMeta.class); + final EnrichedEvent copiedEvent = + rawBarrier.shallowCopySelfAndBindPipeTaskMetaForProgressReport( + "test_pipe", + 1L, + taskMeta, + null, + null, + null, + null, + null, + true, + Long.MIN_VALUE, + Long.MAX_VALUE); + final PipeHeartbeatEvent boundBarrier = (PipeHeartbeatEvent) copiedEvent; + + assertTrue(boundBarrier.isCompletionBarrier()); + assertFalse(boundBarrier.needToCommit()); + assertTrue(boundBarrier.getOnCommittedHooks().isEmpty()); + boundBarrier.bindCompletionSource(4L); + assertTrue(boundBarrier.needToCommit()); + assertEquals(2L, boundBarrier.getAssignerEpoch()); + assertEquals(3L, boundBarrier.getDataGeneration()); + assertEquals(4L, boundBarrier.getCompletionSourceId()); + assertEquals(1, boundBarrier.getOnCommittedHooks().size()); + + final PipeHeartbeatEvent periodicHeartbeat = new PipeHeartbeatEvent(1, false); + final EnrichedEvent boundPeriodicHeartbeat = + periodicHeartbeat.shallowCopySelfAndBindPipeTaskMetaForProgressReport( + "test_pipe", + 1L, + taskMeta, + null, + null, + null, + null, + null, + true, + Long.MIN_VALUE, + Long.MAX_VALUE); + assertFalse(boundPeriodicHeartbeat.needToCommit()); + assertTrue(boundPeriodicHeartbeat.getOnCommittedHooks().isEmpty()); + } +} diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeCompletionOperatorTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeCompletionOperatorTest.java new file mode 100644 index 0000000000000..5a3f9379cde39 --- /dev/null +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeCompletionOperatorTest.java @@ -0,0 +1,252 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.iotdb.db.pipe.metric.overview; + +import org.apache.iotdb.commons.consensus.index.impl.MinimumProgressIndex; +import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta; +import org.apache.iotdb.commons.pipe.agent.task.progress.CommitterKey; +import org.apache.iotdb.db.pipe.event.realtime.PipeRealtimeEvent; +import org.apache.iotdb.db.pipe.source.dataregion.realtime.PipeRealtimeDataRegionSource; +import org.apache.iotdb.db.pipe.source.dataregion.realtime.assigner.PipeDataRegionAssigner; +import org.apache.iotdb.pipe.api.event.Event; + +import org.apache.tsfile.utils.Pair; +import org.junit.Test; + +import java.util.Collections; +import java.util.HashMap; +import java.util.Map; +import java.util.concurrent.atomic.AtomicBoolean; +import java.util.concurrent.atomic.AtomicLong; +import java.util.concurrent.atomic.AtomicReference; + +import static org.junit.Assert.assertEquals; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.when; + +public class PipeDataNodeCompletionOperatorTest { + + @Test + public void testCompletionFollowsCommittedGenerationAndFailsClosed() { + final AtomicBoolean remainingZero = new AtomicBoolean(true); + final AtomicBoolean supported = new AtomicBoolean(true); + final AtomicReference> expected = + new AtomicReference<>(new HashMap<>()); + final AtomicReference currentCommitter = new AtomicReference<>(); + final PipeDataNodeCompletionOperator operator = + new PipeDataNodeCompletionOperator( + remainingZero::get, + () -> new Pair<>(supported.get(), expected.get()), + key -> key == currentCommitter.get()); + + final PipeTaskMeta taskMeta = new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 0); + final PipeTaskMeta staleTaskMeta = new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 0); + final TestSource source = new TestSource(1, taskMeta); + final PipeDataRegionAssigner assigner = mock(PipeDataRegionAssigner.class); + final CommitterKey committerKey = mock(CommitterKey.class); + final CommitterKey nextCommitterKey = mock(CommitterKey.class); + final AtomicLong publishedGeneration = new AtomicLong(0); + final AtomicLong publicationFailureEpoch = new AtomicLong(0); + + when(assigner.getAssignerEpoch()).thenReturn(10L); + when(assigner.getPublishedDataGeneration()).thenAnswer(i -> publishedGeneration.get()); + when(assigner.getPublicationFailureEpoch()).thenAnswer(i -> publicationFailureEpoch.get()); + currentCommitter.set(committerKey); + expected.set(Collections.singletonMap(1, source)); + + operator.registerDataRegionSource(source); + operator.register(source, assigner); + assertEquals(0, operator.getCompletion()); + + operator.markCompleted(1, staleTaskMeta, 10L, 0, source.getCompletionSourceId(), committerKey); + assertEquals(0, operator.getCompletion()); + + operator.markCompleted(1, taskMeta, 10L, 0, source.getCompletionSourceId(), committerKey); + assertEquals(1, operator.getCompletion()); + + publishedGeneration.incrementAndGet(); + assertEquals(0, operator.getCompletion()); + operator.markCompleted(1, taskMeta, 10L, 1, source.getCompletionSourceId(), committerKey); + assertEquals(1, operator.getCompletion()); + + source.setDegraded(true); + assertEquals(0, operator.getCompletion()); + source.setDegraded(false); + assertEquals(1, operator.getCompletion()); + + remainingZero.set(false); + assertEquals(0, operator.getCompletion()); + remainingZero.set(true); + supported.set(false); + assertEquals(0, operator.getCompletion()); + supported.set(true); + assertEquals(1, operator.getCompletion()); + + currentCommitter.set(nextCommitterKey); + assertEquals(0, operator.getCompletion()); + operator.markCompleted(1, taskMeta, 10L, 1, source.getCompletionSourceId(), nextCommitterKey); + assertEquals(1, operator.getCompletion()); + + publicationFailureEpoch.incrementAndGet(); + assertEquals(0, operator.getCompletion()); + + operator.markInvalid(1, taskMeta, source.getCompletionSourceId()); + operator.markCompleted( + 1, taskMeta, 10L, Long.MAX_VALUE, source.getCompletionSourceId(), nextCommitterKey); + assertEquals(0, operator.getCompletion()); + + operator.deregister(source, assigner); + assertEquals(0, operator.getCompletion()); + } + + @Test + public void testMembershipMustMatchExactlyAndEmptyLocalMembershipIsNeutral() { + final AtomicReference> expected = + new AtomicReference<>(new HashMap<>()); + final PipeDataNodeCompletionOperator operator = + new PipeDataNodeCompletionOperator( + () -> true, () -> new Pair<>(true, expected.get()), key -> true); + + final PipeTaskMeta firstTaskMeta = new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 0); + final PipeTaskMeta secondTaskMeta = new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 0); + final TestSource firstSource = new TestSource(1, firstTaskMeta); + final TestSource secondSource = new TestSource(2, secondTaskMeta); + final PipeDataRegionAssigner firstAssigner = mock(PipeDataRegionAssigner.class); + final PipeDataRegionAssigner secondAssigner = mock(PipeDataRegionAssigner.class); + final CommitterKey firstCommitter = mock(CommitterKey.class); + final CommitterKey secondCommitter = mock(CommitterKey.class); + + when(firstAssigner.getAssignerEpoch()).thenReturn(11L); + when(secondAssigner.getAssignerEpoch()).thenReturn(12L); + + operator.registerDataRegionSource(firstSource); + operator.register(firstSource, firstAssigner); + operator.markCompleted( + 1, firstTaskMeta, 11L, 0, firstSource.getCompletionSourceId(), firstCommitter); + expected.set(Collections.singletonMap(1, firstSource)); + assertEquals(1, operator.getCompletion()); + + operator.registerDataRegionSource(secondSource); + operator.register(secondSource, secondAssigner); + operator.markCompleted( + 2, secondTaskMeta, 12L, 0, secondSource.getCompletionSourceId(), secondCommitter); + final Map bothExpected = new HashMap<>(); + bothExpected.put(1, firstSource); + bothExpected.put(2, secondSource); + expected.set(bothExpected); + assertEquals(1, operator.getCompletion()); + + expected.set(Collections.singletonMap(2, secondSource)); + assertEquals(0, operator.getCompletion()); + operator.deregister(firstSource, firstAssigner); + operator.deregisterDataRegionSource(firstSource); + assertEquals(1, operator.getCompletion()); + + expected.set(Collections.emptyMap()); + assertEquals(0, operator.getCompletion()); + operator.deregister(secondSource, secondAssigner); + operator.deregisterDataRegionSource(secondSource); + assertEquals(1, operator.getCompletion()); + } + + @Test + public void testStaleBarrierCannotCompleteReplacementSource() { + final PipeTaskMeta sharedTaskMeta = new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 0); + final PipeDataRegionAssigner sharedAssigner = mock(PipeDataRegionAssigner.class); + final CommitterKey committerKey = mock(CommitterKey.class); + when(sharedAssigner.getAssignerEpoch()).thenReturn(20L); + + final AtomicReference> expected = + new AtomicReference<>(new HashMap<>()); + final PipeDataNodeCompletionOperator operator = + new PipeDataNodeCompletionOperator( + () -> true, () -> new Pair<>(true, expected.get()), key -> true); + + final TestSource oldSource = new TestSource(1, sharedTaskMeta); + operator.registerDataRegionSource(oldSource); + operator.register(oldSource, sharedAssigner); + operator.deregister(oldSource, sharedAssigner); + operator.deregisterDataRegionSource(oldSource); + + final TestSource newSource = new TestSource(1, sharedTaskMeta); + expected.set(Collections.singletonMap(1, newSource)); + operator.registerDataRegionSource(newSource); + operator.register(newSource, sharedAssigner); + + operator.markCompleted( + 1, sharedTaskMeta, 20L, 0, oldSource.getCompletionSourceId(), committerKey); + assertEquals(0, operator.getCompletion()); + + operator.markCompleted( + 1, sharedTaskMeta, 20L, 0, newSource.getCompletionSourceId(), committerKey); + assertEquals(1, operator.getCompletion()); + + operator.markInvalid(1, sharedTaskMeta, oldSource.getCompletionSourceId()); + assertEquals(1, operator.getCompletion()); + + operator.deregister(oldSource, sharedAssigner); + operator.deregisterDataRegionSource(oldSource); + assertEquals(1, operator.getCompletion()); + + operator.markInvalid(1, sharedTaskMeta, newSource.getCompletionSourceId()); + assertEquals(0, operator.getCompletion()); + } + + private static class TestSource extends PipeRealtimeDataRegionSource { + + private final AtomicBoolean degraded = new AtomicBoolean(false); + + private TestSource(final int dataRegionId, final PipeTaskMeta pipeTaskMeta) { + this.dataRegionId = dataRegionId; + this.pipeTaskMeta = pipeTaskMeta; + } + + private void setDegraded(final boolean degraded) { + if (this.degraded.getAndSet(degraded) != degraded) { + markCompletionStateChanged(); + } + } + + @Override + protected void doExtract(final PipeRealtimeEvent event) { + // Do nothing. + } + + @Override + public Event supply() { + return null; + } + + @Override + public boolean isTsFileEpochDegraded() { + return degraded.get(); + } + + @Override + public boolean isNeedListenToTsFile() { + return false; + } + + @Override + public boolean isNeedListenToInsertNode() { + return false; + } + } +} diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeSinglePipeMetricsTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeSinglePipeMetricsTest.java new file mode 100644 index 0000000000000..9368d09a6d505 --- /dev/null +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/metric/overview/PipeDataNodeSinglePipeMetricsTest.java @@ -0,0 +1,221 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.iotdb.db.pipe.metric.overview; + +import org.apache.iotdb.commons.consensus.index.impl.MinimumProgressIndex; +import org.apache.iotdb.commons.pipe.agent.task.meta.PipeTaskMeta; +import org.apache.iotdb.commons.pipe.agent.task.progress.CommitterKey; +import org.apache.iotdb.commons.service.metric.enums.Metric; +import org.apache.iotdb.commons.service.metric.enums.Tag; +import org.apache.iotdb.db.pipe.event.realtime.PipeRealtimeEvent; +import org.apache.iotdb.db.pipe.source.dataregion.IoTDBDataRegionSource; +import org.apache.iotdb.db.pipe.source.dataregion.realtime.PipeRealtimeDataRegionSource; +import org.apache.iotdb.db.pipe.source.dataregion.realtime.assigner.PipeDataRegionAssigner; +import org.apache.iotdb.metrics.AbstractMetricService; +import org.apache.iotdb.metrics.utils.MetricLevel; +import org.apache.iotdb.metrics.utils.MetricType; +import org.apache.iotdb.pipe.api.event.Event; + +import org.junit.Test; +import org.mockito.invocation.Invocation; + +import java.lang.reflect.Field; +import java.util.Arrays; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.mockingDetails; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +public class PipeDataNodeSinglePipeMetricsTest { + + @Test + public void testCompletionGaugeLifecycleAndStaleCallbacksDoNotRecreateState() throws Exception { + final String pipeName = "completion_metric_" + System.nanoTime(); + final long creationTime = 1L; + final String pipeId = pipeName + "_" + creationTime; + final TestSource realtimeSource = + new TestSource( + pipeName, creationTime, 1, new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 0)); + final IoTDBDataRegionSource outerSource = mock(IoTDBDataRegionSource.class); + final PipeDataRegionAssigner assigner = mock(PipeDataRegionAssigner.class); + final AbstractMetricService metricService = mock(AbstractMetricService.class); + + when(outerSource.getPipeName()).thenReturn(pipeName); + when(outerSource.getCreationTime()).thenReturn(creationTime); + when(outerSource.getRealtimeSourceForCompletion()).thenReturn(realtimeSource); + + final PipeDataNodeSinglePipeMetrics metrics = PipeDataNodeSinglePipeMetrics.getInstance(); + final Field metricServiceField = + PipeDataNodeSinglePipeMetrics.class.getDeclaredField("metricService"); + metricServiceField.setAccessible(true); + reset(metrics, metricServiceField); + + try { + metrics.register(outerSource); + metrics.register(realtimeSource, assigner); + metrics.bindTo(metricService); + + assertTrue(hasCompletionGaugeCreation(metricService, pipeName, creationTime)); + + metrics.deregister(pipeId); + assertFalse(metrics.remainingEventAndTimeOperatorMap.containsKey(pipeId)); + assertFalse(metrics.completionOperatorMap.containsKey(pipeId)); + verify(metricService) + .remove( + MetricType.AUTO_GAUGE, + Metric.PIPE_DATANODE_COMPLETION_READY.toString(), + Tag.NAME.toString(), + pipeName, + Tag.CREATION_TIME.toString(), + String.valueOf(creationTime)); + + metrics.register(realtimeSource, assigner); + metrics.markDataRegionCompleted( + pipeName, + creationTime, + 1, + realtimeSource.getPipeTaskMeta(), + 1, + 1, + realtimeSource.getCompletionSourceId(), + mock(CommitterKey.class)); + metrics.markDataRegionInvalid( + pipeName, + creationTime, + 1, + realtimeSource.getPipeTaskMeta(), + realtimeSource.getCompletionSourceId()); + metrics.deregister(realtimeSource, assigner); + metrics.decreaseHeartbeatEventCount(pipeName, creationTime); + assertFalse(metrics.remainingEventAndTimeOperatorMap.containsKey(pipeId)); + assertFalse(metrics.completionOperatorMap.containsKey(pipeId)); + } finally { + reset(metrics, metricServiceField); + } + } + + @Test + public void testDeregisterClearsOperatorsWhenMetricsAreUnbound() throws Exception { + final String pipeName = "completion_unbound_" + System.nanoTime(); + final long creationTime = 2L; + final String pipeId = pipeName + "_" + creationTime; + final TestSource realtimeSource = + new TestSource( + pipeName, creationTime, 2, new PipeTaskMeta(MinimumProgressIndex.INSTANCE, 0)); + final IoTDBDataRegionSource outerSource = mock(IoTDBDataRegionSource.class); + when(outerSource.getPipeName()).thenReturn(pipeName); + when(outerSource.getCreationTime()).thenReturn(creationTime); + when(outerSource.getRealtimeSourceForCompletion()).thenReturn(realtimeSource); + when(outerSource.getHistoricalTsFileInsertionEventCount()).thenReturn(1); + + final PipeDataNodeSinglePipeMetrics metrics = PipeDataNodeSinglePipeMetrics.getInstance(); + final Field metricServiceField = + PipeDataNodeSinglePipeMetrics.class.getDeclaredField("metricService"); + metricServiceField.setAccessible(true); + reset(metrics, metricServiceField); + + try { + metrics.register(outerSource); + assertTrue(metrics.remainingEventAndTimeOperatorMap.containsKey(pipeId)); + assertTrue(metrics.completionOperatorMap.containsKey(pipeId)); + assertEquals( + 1, metrics.remainingEventAndTimeOperatorMap.get(pipeId).getRemainingNonHeartbeatEvents()); + + metrics.deregister(outerSource); + assertEquals( + 0, metrics.remainingEventAndTimeOperatorMap.get(pipeId).getRemainingNonHeartbeatEvents()); + + metrics.deregister(pipeId); + assertFalse(metrics.remainingEventAndTimeOperatorMap.containsKey(pipeId)); + assertFalse(metrics.completionOperatorMap.containsKey(pipeId)); + } finally { + reset(metrics, metricServiceField); + } + } + + private static boolean hasCompletionGaugeCreation( + final AbstractMetricService metricService, final String pipeName, final long creationTime) { + for (final Invocation invocation : mockingDetails(metricService).getInvocations()) { + final Object[] arguments = invocation.getRawArguments(); + if ("createAutoGauge".equals(invocation.getMethod().getName()) + && arguments.length == 5 + && Metric.PIPE_DATANODE_COMPLETION_READY.toString().equals(arguments[0]) + && MetricLevel.IMPORTANT.equals(arguments[1]) + && arguments[2] instanceof PipeDataNodeCompletionOperator + && Arrays.equals( + new String[] { + Tag.NAME.toString(), + pipeName, + Tag.CREATION_TIME.toString(), + String.valueOf(creationTime) + }, + (String[]) arguments[4])) { + return true; + } + } + return false; + } + + private static void reset( + final PipeDataNodeSinglePipeMetrics metrics, final Field metricServiceField) + throws IllegalAccessException { + metrics.remainingEventAndTimeOperatorMap.clear(); + metrics.completionOperatorMap.clear(); + metricServiceField.set(metrics, null); + } + + private static class TestSource extends PipeRealtimeDataRegionSource { + + private TestSource( + final String pipeName, + final long creationTime, + final int dataRegionId, + final PipeTaskMeta pipeTaskMeta) { + this.pipeName = pipeName; + this.creationTime = creationTime; + this.dataRegionId = dataRegionId; + this.pipeTaskMeta = pipeTaskMeta; + } + + @Override + protected void doExtract(final PipeRealtimeEvent event) { + // Do nothing. + } + + @Override + public Event supply() { + return null; + } + + @Override + public boolean isNeedListenToTsFile() { + return false; + } + + @Override + public boolean isNeedListenToInsertNode() { + return false; + } + } +} diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionSourceTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionSourceTest.java index 4677f2126b6c6..d935821cdba4e 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionSourceTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/PipeRealtimeDataRegionSourceTest.java @@ -68,6 +68,25 @@ public void mergedProgressReportEventReleasesNewEvent() throws Exception { } } + @Test + public void completionBarrierIsNeverDroppedByProgressOrHeartbeatMerging() throws Exception { + try (final ProgressReportTestSource source = new ProgressReportTestSource()) { + final PipeRealtimeEvent heartbeatEvent = createHeartbeatEvent(); + final PipeRealtimeEvent completionBarrier = createCompletionBarrierEvent(); + source.extract(heartbeatEvent); + source.extract(completionBarrier); + + Assert.assertFalse(heartbeatEvent.getEvent().isReleased()); + Assert.assertFalse(completionBarrier.getEvent().isReleased()); + Assert.assertEquals(2, source.getEventCount()); + + final PipeRealtimeEvent progressReportEvent = createProgressReportEvent(); + source.extract(progressReportEvent); + Assert.assertFalse(completionBarrier.getEvent().isReleased()); + Assert.assertEquals(3, source.getEventCount()); + } + } + private static PipeRealtimeEvent createHeartbeatEvent() { final PipeRealtimeEvent event = PipeRealtimeEventFactory.createRealtimeEvent(1, false); Assert.assertTrue(event.increaseReferenceCount(TEST_REFERENCE_HOLDER)); @@ -84,6 +103,12 @@ private static PipeRealtimeEvent createProgressReportEvent() { return event; } + private static PipeRealtimeEvent createCompletionBarrierEvent() { + final PipeRealtimeEvent event = PipeRealtimeEventFactory.createCompletionBarrierEvent(1); + Assert.assertTrue(event.increaseReferenceCount(TEST_REFERENCE_HOLDER)); + return event; + } + private static class ProgressReportTestSource extends PipeRealtimeDataRegionSource { @Override diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/assigner/PipeDataRegionAssignerTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/assigner/PipeDataRegionAssignerTest.java new file mode 100644 index 0000000000000..0d58bef247bc0 --- /dev/null +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/assigner/PipeDataRegionAssignerTest.java @@ -0,0 +1,75 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.iotdb.db.pipe.source.dataregion.realtime.assigner; + +import org.apache.iotdb.db.pipe.event.realtime.PipeRealtimeEventFactory; +import org.apache.iotdb.db.pipe.source.dataregion.realtime.assigner.PipeDataRegionAssigner.CompletionToken; + +import org.junit.Test; + +import static org.junit.Assert.assertEquals; +import static org.junit.Assert.assertFalse; +import static org.junit.Assert.assertTrue; + +public class PipeDataRegionAssignerTest { + + @Test + public void testNullDataEventMarksPublicationFailed() { + try (final PipeDataRegionAssigner assigner = new PipeDataRegionAssigner(Integer.MIN_VALUE)) { + final long generation = assigner.getPublishedDataGeneration(); + final long failureEpoch = assigner.getPublicationFailureEpoch(); + + assigner.publishDataEventToAssign(() -> null); + + assertEquals(generation + 1, assigner.getPublishedDataGeneration()); + assertEquals(failureEpoch + 1, assigner.getPublicationFailureEpoch()); + } + } + + @Test + public void testOnlyLatestCompletionInvalidationCanPublishBarrier() { + try (final PipeDataRegionAssigner assigner = new PipeDataRegionAssigner(Integer.MIN_VALUE)) { + final CompletionToken staleToken = assigner.invalidateCompletionAndGetToken(); + final CompletionToken currentToken = assigner.invalidateCompletionAndGetToken(); + + assertFalse(assigner.publishCompletionBarrier(staleToken)); + assigner.invalidateCompletion(); + assertTrue(assigner.publishCompletionBarrier(currentToken)); + } + } + + @Test + public void testInsertInvalidationRejectsCurrentCompletionToken() { + try (final PipeDataRegionAssigner assigner = new PipeDataRegionAssigner(Integer.MIN_VALUE)) { + final long generation = assigner.getPublishedDataGeneration(); + final CompletionToken ignoredInsertToken = assigner.invalidateCompletionAndGetToken(); + + assigner.invalidateCompletionBarrier(); + + assertEquals(generation + 2, assigner.getPublishedDataGeneration()); + assertFalse(assigner.publishCompletionBarrier(ignoredInsertToken)); + + final CompletionToken publishedInsertToken = assigner.invalidateCompletionAndGetToken(); + assigner.publishInsertDataEventToAssign( + () -> PipeRealtimeEventFactory.createRealtimeEvent(Integer.MIN_VALUE, false)); + assertFalse(assigner.publishCompletionBarrier(publishedInsertToken)); + } + } +} diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/listener/PipeInsertionDataNodeListenerTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/listener/PipeInsertionDataNodeListenerTest.java new file mode 100644 index 0000000000000..08240ae24f76a --- /dev/null +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/pipe/source/dataregion/realtime/listener/PipeInsertionDataNodeListenerTest.java @@ -0,0 +1,66 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one + * or more contributor license agreements. See the NOTICE file + * distributed with this work for additional information + * regarding copyright ownership. The ASF licenses this file + * to you under the Apache License, Version 2.0 (the + * "License"); you may not use this file except in compliance + * with the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, + * software distributed under the License is distributed on an + * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY + * KIND, either express or implied. See the License for the + * specific language governing permissions and limitations + * under the License. + */ + +package org.apache.iotdb.db.pipe.source.dataregion.realtime.listener; + +import org.apache.iotdb.db.pipe.source.dataregion.realtime.assigner.PipeDataRegionAssigner; + +import org.junit.Test; + +import java.lang.reflect.Field; +import java.util.concurrent.ConcurrentMap; + +import static org.mockito.ArgumentMatchers.any; +import static org.mockito.Mockito.mock; +import static org.mockito.Mockito.never; +import static org.mockito.Mockito.reset; +import static org.mockito.Mockito.verify; +import static org.mockito.Mockito.when; + +public class PipeInsertionDataNodeListenerTest { + + @Test + @SuppressWarnings("unchecked") + public void testInsertUsesBarrierInvalidatingPublicationPath() throws Exception { + final PipeInsertionDataNodeListener listener = PipeInsertionDataNodeListener.getInstance(); + final Field assignerMapField = + PipeInsertionDataNodeListener.class.getDeclaredField("dataRegionId2Assigner"); + assignerMapField.setAccessible(true); + final ConcurrentMap assignerMap = + (ConcurrentMap) assignerMapField.get(listener); + final int dataRegionId = Integer.MIN_VALUE + 1; + final PipeDataRegionAssigner assigner = mock(PipeDataRegionAssigner.class); + assignerMap.put(dataRegionId, assigner); + + try { + when(assigner.shouldListenToInsertNode()).thenReturn(false); + listener.listenToInsertNode(dataRegionId, null, null, null); + verify(assigner).invalidateCompletionBarrier(); + verify(assigner, never()).invalidateCompletion(); + + reset(assigner); + when(assigner.shouldListenToInsertNode()).thenReturn(true); + listener.listenToInsertNode(dataRegionId, null, null, null); + verify(assigner).publishInsertDataEventToAssign(any()); + verify(assigner, never()).publishDataEventToAssign(any()); + } finally { + assignerMap.remove(dataRegionId, assigner); + } + } +} diff --git a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/DataRegionTest.java b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/DataRegionTest.java index 1000ad31cc4c0..b668a24c14587 100644 --- a/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/DataRegionTest.java +++ b/iotdb-core/datanode/src/test/java/org/apache/iotdb/db/storageengine/dataregion/DataRegionTest.java @@ -65,8 +65,10 @@ import org.apache.iotdb.db.storageengine.dataregion.compaction.selector.constant.InnerSequenceCompactionSelector; import org.apache.iotdb.db.storageengine.dataregion.compaction.selector.constant.InnerUnsequenceCompactionSelector; import org.apache.iotdb.db.storageengine.dataregion.compaction.utils.CompactionConfigRestorer; +import org.apache.iotdb.db.storageengine.dataregion.flush.FlushListener; import org.apache.iotdb.db.storageengine.dataregion.flush.FlushManager; import org.apache.iotdb.db.storageengine.dataregion.flush.TsFileFlushPolicy; +import org.apache.iotdb.db.storageengine.dataregion.memtable.IMemTable; import org.apache.iotdb.db.storageengine.dataregion.memtable.ReadOnlyMemChunk; import org.apache.iotdb.db.storageengine.dataregion.memtable.TsFileProcessor; import org.apache.iotdb.db.storageengine.dataregion.modification.DeletionPredicate; @@ -112,8 +114,12 @@ import java.util.List; import java.util.Map; import java.util.concurrent.CompletableFuture; +import java.util.concurrent.CompletionException; +import java.util.concurrent.CountDownLatch; import java.util.concurrent.ExecutionException; import java.util.concurrent.Future; +import java.util.concurrent.TimeUnit; +import java.util.concurrent.TimeoutException; import java.util.concurrent.atomic.AtomicInteger; import static org.apache.iotdb.db.queryengine.plan.statement.StatementTestUtils.genInsertRowNode; @@ -2262,4 +2268,141 @@ public void testFlushSpecifiedResource() future.get(); assertTrue(tsFileResourceSeq.isClosed()); } + + @Test + public void testFullFlushClosesProcessorAfterOrdinaryAsyncFlush() throws Exception { + final TSRecord record = new TSRecord(deviceId, 100); + record.addTuple(DataPoint.getDataPoint(TSDataType.INT32, measurementId, String.valueOf(100))); + dataRegion.insert(buildInsertRowNodeByTSRecord(record)); + + final TsFileProcessor tsFileProcessor = + dataRegion.getWorkSequenceTsFileProcessors().iterator().next(); + final TsFileResource tsFileResource = tsFileProcessor.getTsFileResource(); + tsFileProcessor.asyncFlush(); + final Future ordinaryFlushFuture = tsFileProcessor.getCloseFuture(); + Assert.assertNotNull(ordinaryFlushFuture); + ordinaryFlushFuture.get(); + + Assert.assertFalse(tsFileProcessor.alreadyMarkedClosing()); + Assert.assertFalse(tsFileResource.isClosed()); + + dataRegion.syncCloseAllWorkingAndClosingTsFileProcessors(); + + Assert.assertTrue(tsFileResource.isClosed()); + Assert.assertFalse(dataRegion.getWorkSequenceTsFileProcessors().contains(tsFileProcessor)); + } + + @Test + public void testFullFlushWaitsForRunningOrdinaryAsyncFlush() throws Exception { + final CountDownLatch ordinaryFlushBlocked = new CountDownLatch(1); + final CountDownLatch allowOrdinaryFlush = new CountDownLatch(1); + dataRegion.setCustomFlushListeners( + Collections.singletonList( + new FlushListener() { + @Override + public void onMemTableFlushStarted(final IMemTable memTable) { + // Do nothing. + } + + @Override + public void onMemTableFlushed(final IMemTable memTable) { + ordinaryFlushBlocked.countDown(); + try { + allowOrdinaryFlush.await(); + } catch (final InterruptedException e) { + Thread.currentThread().interrupt(); + throw new AssertionError(e); + } + } + })); + + final TSRecord record = new TSRecord(deviceId, 100); + record.addTuple(DataPoint.getDataPoint(TSDataType.INT32, measurementId, String.valueOf(100))); + dataRegion.insert(buildInsertRowNodeByTSRecord(record)); + final TsFileProcessor tsFileProcessor = + dataRegion.getWorkSequenceTsFileProcessors().iterator().next(); + final TsFileResource tsFileResource = tsFileProcessor.getTsFileResource(); + + tsFileProcessor.asyncFlush(); + Assert.assertTrue(ordinaryFlushBlocked.await(30, TimeUnit.SECONDS)); + final CompletableFuture fullFlushFuture = + CompletableFuture.runAsync( + () -> { + try { + dataRegion.syncCloseAllWorkingAndClosingTsFileProcessors(); + } catch (final InterruptedException e) { + Thread.currentThread().interrupt(); + throw new CompletionException(e); + } catch (final ExecutionException e) { + throw new CompletionException(e); + } + }); + + try { + fullFlushFuture.get(200, TimeUnit.MILLISECONDS); + Assert.fail(); + } catch (final TimeoutException expected) { + // The full flush must wait until the ordinary flush releases the processor. + } finally { + allowOrdinaryFlush.countDown(); + } + + fullFlushFuture.get(30, TimeUnit.SECONDS); + Assert.assertTrue(tsFileResource.isClosed()); + Assert.assertFalse(dataRegion.getWorkSequenceTsFileProcessors().contains(tsFileProcessor)); + } + + @Test + public void testFullFlushDoesNotChaseProcessorCreatedAfterSnapshot() throws Exception { + final CountDownLatch closeListenerBlocked = new CountDownLatch(1); + final CountDownLatch allowCloseListener = new CountDownLatch(1); + dataRegion.setCustomCloseFileListeners( + Collections.singletonList( + processor -> { + closeListenerBlocked.countDown(); + try { + allowCloseListener.await(); + } catch (final InterruptedException e) { + Thread.currentThread().interrupt(); + throw new TsFileProcessorException(e); + } + })); + + final TSRecord firstRecord = new TSRecord(deviceId, 100); + firstRecord.addTuple( + DataPoint.getDataPoint(TSDataType.INT32, measurementId, String.valueOf(100))); + dataRegion.insert(buildInsertRowNodeByTSRecord(firstRecord)); + final TsFileProcessor initialProcessor = + dataRegion.getWorkSequenceTsFileProcessors().iterator().next(); + + final CompletableFuture fullFlushFuture = + CompletableFuture.runAsync( + () -> { + try { + dataRegion.syncCloseAllWorkingAndClosingTsFileProcessors(); + } catch (final InterruptedException e) { + Thread.currentThread().interrupt(); + throw new CompletionException(e); + } catch (final ExecutionException e) { + throw new CompletionException(e); + } + }); + + TsFileProcessor replacementProcessor; + try { + Assert.assertTrue(closeListenerBlocked.await(30, TimeUnit.SECONDS)); + final TSRecord secondRecord = new TSRecord(deviceId, 101); + secondRecord.addTuple( + DataPoint.getDataPoint(TSDataType.INT32, measurementId, String.valueOf(101))); + dataRegion.insert(buildInsertRowNodeByTSRecord(secondRecord)); + replacementProcessor = dataRegion.getWorkSequenceTsFileProcessors().iterator().next(); + Assert.assertNotSame(initialProcessor, replacementProcessor); + } finally { + allowCloseListener.countDown(); + } + + fullFlushFuture.get(30, TimeUnit.SECONDS); + Assert.assertTrue(dataRegion.getWorkSequenceTsFileProcessors().contains(replacementProcessor)); + Assert.assertFalse(replacementProcessor.alreadyMarkedClosing()); + } } diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/meta/PipeTaskMeta.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/meta/PipeTaskMeta.java index e9939d7b2c6f5..3873ae090918c 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/meta/PipeTaskMeta.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/meta/PipeTaskMeta.java @@ -39,6 +39,7 @@ import java.util.Set; import java.util.concurrent.ConcurrentHashMap; import java.util.concurrent.atomic.AtomicInteger; +import java.util.concurrent.atomic.AtomicLong; import java.util.concurrent.atomic.AtomicReference; public class PipeTaskMeta { @@ -60,6 +61,8 @@ public class PipeTaskMeta { private final Set exceptionMessages = Collections.newSetFromMap(new ConcurrentHashMap<>()); + private final AtomicLong exceptionMessageVersion = new AtomicLong(0); + public PipeTaskMeta(/* @NotNull */ final ProgressIndex progressIndex, final int leaderNodeId) { this.progressIndex.set(progressIndex); this.leaderNodeId.set(leaderNodeId); @@ -119,6 +122,7 @@ public synchronized void trackExceptionMessage(final PipeRuntimeException except // Here we still keep the map form to allow compatibility with legacy versions exceptionMessages.clear(); exceptionMessages.add(exceptionMessage); + exceptionMessageVersion.incrementAndGet(); } public synchronized boolean containsExceptionMessage( @@ -131,11 +135,20 @@ public synchronized boolean hasExceptionMessages() { } public synchronized void clearExceptionMessages() { - exceptionMessages.clear(); + if (!exceptionMessages.isEmpty()) { + exceptionMessages.clear(); + exceptionMessageVersion.incrementAndGet(); + } } public synchronized void clearExceptionMessagesBefore(final long exceptionsClearTime) { - exceptionMessages.removeIf(exception -> exception.getTimeStamp() <= exceptionsClearTime); + if (exceptionMessages.removeIf(exception -> exception.getTimeStamp() <= exceptionsClearTime)) { + exceptionMessageVersion.incrementAndGet(); + } + } + + public long getExceptionMessageVersion() { + return exceptionMessageVersion.get(); } public synchronized void serialize(final OutputStream outputStream) throws IOException { diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/progress/PipeEventCommitManager.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/progress/PipeEventCommitManager.java index 26e7ea305d595..b21861ad67114 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/progress/PipeEventCommitManager.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/pipe/agent/task/progress/PipeEventCommitManager.java @@ -173,6 +173,16 @@ public CommitterKey getCommitterKey( return generateCommitterKey(pipeName, creationTime, regionId); } + public boolean isCurrentCommitterKey(final CommitterKey committerKey) { + return committerKey != null + && eventCommitterMap.containsKey(committerKey) + && committerKey.equals( + generateCommitterKey( + committerKey.getPipeName(), + committerKey.getCreationTime(), + committerKey.getRegionId())); + } + private CommitterKey generateCommitterKey( final String pipeName, final long creationTime, final int regionId) { return taskAgent.getCommitterKey( diff --git a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/enums/Metric.java b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/enums/Metric.java index f5a76e6d66a6a..996af79cbe1d8 100644 --- a/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/enums/Metric.java +++ b/iotdb-core/node-commons/src/main/java/org/apache/iotdb/commons/service/metric/enums/Metric.java @@ -189,6 +189,7 @@ public enum Metric { PIPE_CONNECTOR_SCHEMA_TRANSFER("pipe_connector_schema_transfer"), PIPE_DATANODE_REMAINING_EVENT_COUNT("pipe_datanode_remaining_event_count"), PIPE_DATANODE_REMAINING_TIME("pipe_datanode_remaining_time"), + PIPE_DATANODE_COMPLETION_READY("pipe_datanode_completion_ready"), PIPE_INSERT_NODE_EVENT_TRANSFER_TIME("pipe_insert_node_event_transfer_time"), PIPE_TSFILE_EVENT_TRANSFER_TIME("pipe_tsfile_event_transfer_time"), PIPE_DATANODE_EVENT_TRANSFER("pipe_datanode_event_transfer"),