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 @@ -45,6 +45,7 @@
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 @@ -63,7 +64,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 @@ -75,6 +79,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 @@ -118,6 +123,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());
Comment on lines +130 to +139

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is it possible to support air-gap?

Comment on lines +126 to +139

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

How about adding a field supportCompletionMark for these enums?


protected static final IoTDBConfig CONFIG = IoTDBDescriptor.getInstance().getConfig();

Expand Down Expand Up @@ -677,6 +696,90 @@ 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) {
Comment on lines +699 to +700

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Add a comment to explain the return value.

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);
Comment on lines +726 to +744

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

How about caching the result in pipeMeta?


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;
}
}
}
if (!dataRegionId2Source.keySet().equals(getExpectedDataRegionIds(pipeMeta))) {
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 @@ -134,6 +134,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 @@ -36,6 +36,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 @@ -60,6 +61,8 @@ public class PipeEventCollector implements EventCollector {

private final int regionId;

private final long completionSourceId;

Comment on lines +64 to +65

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is it possible to make the ID not only for completion?

private final boolean forceTabletFormat;

private final boolean skipParsing;
Expand All @@ -76,12 +79,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 @@ -246,6 +251,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 @@ -255,6 +261,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 @@ -275,6 +284,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(
Comment on lines +294 to +295

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Is the method name markDataRegionInvalid proper?

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();
}
Comment on lines +117 to +124

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Why must it be a realtime source

}
Original file line number Diff line number Diff line change
Expand Up @@ -228,7 +228,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);
}
Comment on lines +231 to +237

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

May use a higher-level abstraction like Event.isUserVisible or Event.shouldSkipProcessing

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