Package | Description |
---|---|
org.apache.flink.streaming.api.operators | |
org.apache.flink.streaming.runtime.io | |
org.apache.flink.streaming.runtime.tasks |
This package contains classes that realize streaming tasks.
|
Modifier and Type | Method and Description |
---|---|
void |
StreamSource.run(Object lockingObject,
StreamStatusMaintainer streamStatusMaintainer,
OperatorChain<?,?> operatorChain) |
void |
StreamSource.run(Object lockingObject,
StreamStatusMaintainer streamStatusMaintainer,
Output<StreamRecord<OUT>> collector,
OperatorChain<?,?> operatorChain) |
Modifier and Type | Method and Description |
---|---|
static StreamMultipleInputProcessor |
StreamMultipleInputProcessorFactory.create(AbstractInvokable ownerTask,
CheckpointedInputGate[] checkpointedInputGates,
StreamConfig.InputConfig[] configuredInputs,
IOManager ioManager,
MemoryManager memoryManager,
TaskIOMetricGroup ioMetricGroup,
Counter mainOperatorRecordsIn,
StreamStatusMaintainer streamStatusMaintainer,
MultipleInputStreamOperator<?> mainOperator,
WatermarkGauge[] inputWatermarkGauges,
StreamConfig streamConfig,
Configuration taskManagerConfig,
Configuration jobConfig,
ExecutionConfig executionConfig,
ClassLoader userClassloader,
OperatorChain<?,?> operatorChain) |
Modifier and Type | Field and Description |
---|---|
protected OperatorChain<OUT,OP> |
StreamTask.operatorChain
The chain of operators executed by this task.
|
Modifier and Type | Method and Description |
---|---|
void |
SubtaskCheckpointCoordinator.abortCheckpointOnBarrier(long checkpointId,
Throwable cause,
OperatorChain<?,?> operatorChain) |
void |
SubtaskCheckpointCoordinator.checkpointState(CheckpointMetaData checkpointMetaData,
CheckpointOptions checkpointOptions,
CheckpointMetricsBuilder checkpointMetrics,
OperatorChain<?,?> operatorChain,
java.util.function.Supplier<Boolean> isRunning)
Must be called after
SubtaskCheckpointCoordinator.initCheckpoint(long, CheckpointOptions) . |
void |
SubtaskCheckpointCoordinator.notifyCheckpointAborted(long checkpointId,
OperatorChain<?,?> operatorChain,
java.util.function.Supplier<Boolean> isRunning)
Notified on the task side once a distributed checkpoint has been aborted.
|
void |
SubtaskCheckpointCoordinator.notifyCheckpointComplete(long checkpointId,
OperatorChain<?,?> operatorChain,
java.util.function.Supplier<Boolean> isRunning)
Notified on the task side once a distributed checkpoint has been completed.
|
Copyright © 2014–2021 The Apache Software Foundation. All rights reserved.