Package | Description |
---|---|
org.apache.flink.streaming.api.operators | |
org.apache.flink.streaming.api.operators.sort | |
org.apache.flink.streaming.api.operators.source | |
org.apache.flink.streaming.runtime.io | |
org.apache.flink.streaming.runtime.tasks |
This package contains classes that realize streaming tasks.
|
org.apache.flink.streaming.runtime.watermarkstatus |
Modifier and Type | Method and Description |
---|---|
DataInputStatus |
SourceOperator.emitNext(PushingAsyncDataInput.DataOutput<OUT> output) |
Modifier and Type | Method and Description |
---|---|
DataInputStatus |
MultiInputSortingDataInput.emitNext(PushingAsyncDataInput.DataOutput<IN> output) |
DataInputStatus |
SortingDataInput.emitNext(PushingAsyncDataInput.DataOutput<T> output) |
Modifier and Type | Method and Description |
---|---|
ReaderOutput<T> |
TimestampsAndWatermarks.createMainOutput(PushingAsyncDataInput.DataOutput<T> output,
TimestampsAndWatermarks.WatermarkUpdateListener watermarkCallback)
Creates the ReaderOutput for the source reader, than internally runs the timestamp extraction
and watermark generation.
|
ReaderOutput<T> |
NoOpTimestampsAndWatermarks.createMainOutput(PushingAsyncDataInput.DataOutput<T> output,
TimestampsAndWatermarks.WatermarkUpdateListener watermarkEmitted) |
ReaderOutput<T> |
ProgressiveTimestampsAndWatermarks.createMainOutput(PushingAsyncDataInput.DataOutput<T> output,
TimestampsAndWatermarks.WatermarkUpdateListener watermarkEmitted) |
static <E> SourceOutputWithWatermarks<E> |
SourceOutputWithWatermarks.createWithSeparateOutputs(PushingAsyncDataInput.DataOutput<E> recordsOutput,
WatermarkOutput onEventWatermarkOutput,
WatermarkOutput periodicWatermarkOutput,
TimestampAssigner<E> timestampAssigner,
WatermarkGenerator<E> watermarkGenerator)
Creates a new SourceOutputWithWatermarks that emits records to the given DataOutput and
watermarks to the different WatermarkOutputs.
|
Constructor and Description |
---|
SourceOutputWithWatermarks(PushingAsyncDataInput.DataOutput<T> recordsOutput,
WatermarkOutput onEventWatermarkOutput,
WatermarkOutput periodicWatermarkOutput,
TimestampAssigner<T> timestampAssigner,
WatermarkGenerator<T> watermarkGenerator)
Creates a new SourceOutputWithWatermarks that emits records to the given DataOutput and
watermarks to the (possibly different) WatermarkOutput.
|
WatermarkToDataOutput(PushingAsyncDataInput.DataOutput<?> output) |
WatermarkToDataOutput(PushingAsyncDataInput.DataOutput<?> output,
TimestampsAndWatermarks.WatermarkUpdateListener watermarkEmitted)
Creates a new WatermarkOutput against the given DataOutput.
|
Modifier and Type | Class and Description |
---|---|
class |
FinishedDataOutput<IN>
An empty
PushingAsyncDataInput.DataOutput which is used by StreamOneInputProcessor once an DataInputStatus.END_OF_DATA is received. |
Modifier and Type | Method and Description |
---|---|
DataInputStatus |
StreamTaskSourceInput.emitNext(PushingAsyncDataInput.DataOutput<T> output) |
DataInputStatus |
AbstractStreamTaskNetworkInput.emitNext(PushingAsyncDataInput.DataOutput<T> output) |
DataInputStatus |
StreamTaskExternallyInducedSourceInput.emitNext(PushingAsyncDataInput.DataOutput<T> output) |
DataInputStatus |
PushingAsyncDataInput.emitNext(PushingAsyncDataInput.DataOutput<T> output)
Pushes the next element to the output from current data input, and returns the input status
to indicate whether there are more available data in current input.
|
Constructor and Description |
---|
StreamOneInputProcessor(StreamTaskInput<IN> input,
PushingAsyncDataInput.DataOutput<IN> output,
BoundedMultiInput endOfInputAware) |
Modifier and Type | Class and Description |
---|---|
static class |
SourceOperatorStreamTask.AsyncDataOutputToOutput<T>
Implementation of
PushingAsyncDataInput.DataOutput that wraps a specific Output . |
Modifier and Type | Method and Description |
---|---|
DataInputStatus |
StreamTaskFinishedOnRestoreSourceInput.emitNext(PushingAsyncDataInput.DataOutput<T> output) |
Modifier and Type | Method and Description |
---|---|
void |
StatusWatermarkValve.inputWatermark(Watermark watermark,
int channelIndex,
PushingAsyncDataInput.DataOutput<?> output)
Feed a
Watermark into the valve. |
void |
StatusWatermarkValve.inputWatermarkStatus(WatermarkStatus watermarkStatus,
int channelIndex,
PushingAsyncDataInput.DataOutput<?> output)
Feed a
WatermarkStatus into the valve. |
Copyright © 2014–2024 The Apache Software Foundation. All rights reserved.