Package | Description |
---|---|
org.apache.flink.runtime.io.network.partition.consumer | |
org.apache.flink.runtime.taskmanager | |
org.apache.flink.streaming.runtime.io | |
org.apache.flink.streaming.runtime.io.checkpointing | |
org.apache.flink.streaming.runtime.tasks |
This package contains classes that realize streaming tasks.
|
Modifier and Type | Class and Description |
---|---|
class |
IndexedInputGate
An
InputGate with a specific index. |
class |
SingleInputGate
An input gate consumes one or more partitions of a single produced intermediate result.
|
Modifier and Type | Class and Description |
---|---|
class |
InputGateWithMetrics
This class wraps
InputGate provided by shuffle service and it is mainly used for
increasing general input metrics from TaskIOMetricGroup . |
Modifier and Type | Class and Description |
---|---|
class |
StreamTaskExternallyInducedSourceInput<T>
A subclass of
StreamTaskSourceInput for ExternallyInducedSourceReader . |
class |
StreamTaskSourceInput<T>
Implementation of
StreamTaskInput that reads data from the SourceOperator and
returns the DataInputStatus to indicate whether the source state is available,
unavailable or finished. |
Modifier and Type | Method and Description |
---|---|
static SingleCheckpointBarrierHandler |
SingleCheckpointBarrierHandler.aligned(String taskName,
CheckpointableTask toNotifyOnCheckpoint,
Clock clock,
int numOpenChannels,
java.util.function.BiFunction<Callable<?>,java.time.Duration,org.apache.flink.streaming.runtime.io.checkpointing.CheckpointBarrierHandler.Cancellable> registerTimer,
boolean enableCheckpointAfterTasksFinished,
CheckpointableInput... inputs) |
static SingleCheckpointBarrierHandler |
SingleCheckpointBarrierHandler.alternating(String taskName,
CheckpointableTask toNotifyOnCheckpoint,
SubtaskCheckpointCoordinator checkpointCoordinator,
Clock clock,
int numOpenChannels,
java.util.function.BiFunction<Callable<?>,java.time.Duration,org.apache.flink.streaming.runtime.io.checkpointing.CheckpointBarrierHandler.Cancellable> registerTimer,
boolean enableCheckpointAfterTasksFinished,
CheckpointableInput... inputs) |
static SingleCheckpointBarrierHandler |
SingleCheckpointBarrierHandler.createUnalignedCheckpointBarrierHandler(SubtaskCheckpointCoordinator checkpointCoordinator,
String taskName,
CheckpointableTask toNotifyOnCheckpoint,
Clock clock,
boolean enableCheckpointsAfterTasksFinish,
CheckpointableInput... inputs) |
static SingleCheckpointBarrierHandler |
SingleCheckpointBarrierHandler.unaligned(String taskName,
CheckpointableTask toNotifyOnCheckpoint,
SubtaskCheckpointCoordinator checkpointCoordinator,
Clock clock,
int numOpenChannels,
java.util.function.BiFunction<Callable<?>,java.time.Duration,org.apache.flink.streaming.runtime.io.checkpointing.CheckpointBarrierHandler.Cancellable> registerTimer,
boolean enableCheckpointAfterTasksFinished,
CheckpointableInput... inputs) |
Modifier and Type | Class and Description |
---|---|
class |
StreamTaskFinishedOnRestoreSourceInput<T>
A special source input implementation that immediately emit END_OF_INPUT.
|
Copyright © 2014–2023 The Apache Software Foundation. All rights reserved.