Package | Description |
---|---|
org.apache.flink.streaming.api.operators.sort | |
org.apache.flink.streaming.runtime.io | |
org.apache.flink.streaming.runtime.io.recovery |
Modifier and Type | Class and Description |
---|---|
class |
MultiInputSortingDataInput<IN,K>
An input that wraps an underlying input and sorts the incoming records.
|
class |
SortingDataInput<T,K>
A
StreamTaskInput which sorts in the incoming records from a chained input. |
Modifier and Type | Method and Description |
---|---|
StreamTaskInput<?>[] |
MultiInputSortingDataInput.SelectableSortingInputs.getPassThroughInputs() |
StreamTaskInput<?>[] |
MultiInputSortingDataInput.SelectableSortingInputs.getSortedInputs() |
Modifier and Type | Method and Description |
---|---|
static <K> MultiInputSortingDataInput.SelectableSortingInputs |
MultiInputSortingDataInput.wrapInputs(AbstractInvokable containingTask,
StreamTaskInput<Object>[] sortingInputs,
KeySelector<Object,K>[] keySelectors,
TypeSerializer<Object>[] inputSerializers,
TypeSerializer<K> keySerializer,
StreamTaskInput<Object>[] passThroughInputs,
MemoryManager memoryManager,
IOManager ioManager,
boolean objectReuse,
double managedMemoryFraction,
Configuration jobConfiguration) |
static <K> MultiInputSortingDataInput.SelectableSortingInputs |
MultiInputSortingDataInput.wrapInputs(AbstractInvokable containingTask,
StreamTaskInput<Object>[] sortingInputs,
KeySelector<Object,K>[] keySelectors,
TypeSerializer<Object>[] inputSerializers,
TypeSerializer<K> keySerializer,
StreamTaskInput<Object>[] passThroughInputs,
MemoryManager memoryManager,
IOManager ioManager,
boolean objectReuse,
double managedMemoryFraction,
Configuration jobConfiguration) |
Constructor and Description |
---|
SelectableSortingInputs(StreamTaskInput<?>[] sortedInputs,
StreamTaskInput<?>[] passThroughInputs,
InputSelectable inputSelectable) |
SelectableSortingInputs(StreamTaskInput<?>[] sortedInputs,
StreamTaskInput<?>[] passThroughInputs,
InputSelectable inputSelectable) |
SortingDataInput(StreamTaskInput<T> wrappedInput,
TypeSerializer<T> typeSerializer,
TypeSerializer<K> keySerializer,
KeySelector<T,K> keySelector,
MemoryManager memoryManager,
IOManager ioManager,
boolean objectReuse,
double managedMemoryFraction,
Configuration jobConfiguration,
AbstractInvokable containingTask) |
Modifier and Type | Interface and Description |
---|---|
interface |
RecoverableStreamTaskInput<T>
A
StreamTaskInput used during recovery of in-flight data. |
Modifier and Type | Class and Description |
---|---|
class |
AbstractStreamTaskNetworkInput<T,R extends RecordDeserializer<DeserializationDelegate<StreamElement>>>
Base class for network-based StreamTaskInput where each channel has a designated
RecordDeserializer for spanning records. |
class |
StreamTaskExternallyInducedSourceInput<T>
A subclass of
StreamTaskSourceInput for ExternallyInducedSourceReader . |
class |
StreamTaskNetworkInput<T>
Implementation of
StreamTaskInput that wraps an input from network taken from CheckpointedInputGate . |
class |
StreamTaskSourceInput<T>
Implementation of
StreamTaskInput that reads data from the SourceOperator and
returns the InputStatus to indicate whether the source state is available, unavailable or
finished. |
Modifier and Type | Method and Description |
---|---|
static <T> StreamTaskInput<T> |
StreamTaskNetworkInputFactory.create(CheckpointedInputGate checkpointedInputGate,
TypeSerializer<T> inputSerializer,
IOManager ioManager,
StatusWatermarkValve statusWatermarkValve,
int inputIndex,
InflightDataRescalingDescriptor rescalingDescriptorinflightDataRescalingDescriptor,
java.util.function.Function<Integer,StreamPartitioner<?>> gatePartitioners,
TaskInfo taskInfo)
Factory method for
StreamTaskNetworkInput or RescalingStreamTaskNetworkInput
depending on InflightDataRescalingDescriptor . |
StreamTaskInput<T> |
RecoverableStreamTaskInput.finishRecovery() |
Constructor and Description |
---|
StreamOneInputProcessor(StreamTaskInput<IN> input,
PushingAsyncDataInput.DataOutput<IN> output,
BoundedMultiInput endOfInputAware) |
Modifier and Type | Class and Description |
---|---|
class |
RescalingStreamTaskNetworkInput<T>
A
StreamTaskNetworkInput implementation that demultiplexes virtual channels. |
Modifier and Type | Method and Description |
---|---|
StreamTaskInput<T> |
RescalingStreamTaskNetworkInput.finishRecovery() |
Copyright © 2014–2022 The Apache Software Foundation. All rights reserved.