Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -109,6 +111,10 @@ public String getPipeName() {
return pipeName;
}

public PipeExtractor getPipeExtractor() {
return ((PipeTaskSourceStage) sourceStage).getPipeExtractor();
}

public boolean isCompleted() {
return isCompleted;
}
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand All @@ -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;
Expand Down Expand Up @@ -115,6 +121,20 @@
public class PipeDataNodeTaskAgent extends PipeTaskAgent {

private static final Logger LOGGER = LoggerFactory.getLogger(PipeDataNodeTaskAgent.class);
private static final Set<String> COMPLETION_SUPPORTED_SOURCES =
Set.of(
BuiltinPipePlugin.IOTDB_EXTRACTOR.getPipePluginName(),
BuiltinPipePlugin.IOTDB_SOURCE.getPipePluginName());
private static final Set<String> 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();

Expand Down Expand Up @@ -570,6 +590,87 @@ public Set<Integer> getPipeTaskRegionIdSet(final String pipeName, final long cre
: pipeMeta.getRuntimeMeta().getConsensusGroupId2TaskMetaMap().keySet();
}

public Pair<Boolean, Map<Integer, PipeRealtimeDataRegionSource>> 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<Integer, PipeRealtimeDataRegionSource> dataRegionId2Source = new HashMap<>();
final Map<Integer, PipeTask> pipeTaskMap = pipeTaskManager.getPipeTasks(staticMeta);
if (pipeTaskMap != null) {
for (final Map.Entry<Integer, PipeTask> 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);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -125,6 +125,7 @@ public PipeDataNodeTask build() {
pipeStaticMeta.getCreationTime(),
blendUserAndSystemParameters(pipeStaticMeta.getProcessorParameters(), pipeTaskMeta),
regionId,
sourceStage.getCompletionSourceId(),
sourceStage.getEventSupplier(),
sinkStage.getPipeSinkPendingQueue(),
PROCESSOR_EXECUTOR,
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand All @@ -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;
Expand All @@ -72,12 +75,14 @@ public PipeEventCollector(
final UnboundedBlockingPendingQueue<Event> 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;
Expand Down Expand Up @@ -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;
Expand All @@ -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());
Expand All @@ -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);
}
}

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -67,6 +67,7 @@ public PipeTaskProcessorStage(
final long creationTime,
final PipeParameters pipeProcessorParameters,
final int regionId,
final long completionSourceId,
final EventSupplier pipeSourceInputEventSupplier,
final UnboundedBlockingPendingQueue<Event> pipeSinkOutputPendingQueue,
final PipeProcessorSubtaskExecutor executor,
Expand Down Expand Up @@ -111,6 +112,7 @@ public PipeTaskProcessorStage(
pipeSinkOutputPendingQueue,
creationTime,
regionId,
completionSourceId,
forceTabletFormat,
skipParsing,
isUsedForConsensusPipe);
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down Expand Up @@ -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();
}
}
Original file line number Diff line number Diff line change
Expand Up @@ -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 {
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -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;
Expand Down
Loading
Loading