Modifier and Type | Method and Description |
---|---|
void |
AbstractStreamingWriter.processWatermark(Watermark mark) |
void |
PartitionCommitter.processWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
void |
StateBootstrapWrapperOperator.processWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
Watermark |
AssignerWithPunctuatedWatermarks.checkAndGetNextWatermark(T lastElement,
long extractedTimestamp)
Deprecated.
Asks this implementation if it wants to emit a watermark.
|
Watermark |
AssignerWithPeriodicWatermarks.getCurrentWatermark()
Deprecated.
Returns the current watermark.
|
Watermark |
IngestionTimeExtractor.getCurrentWatermark()
Deprecated.
|
Modifier and Type | Method and Description |
---|---|
void |
SourceFunction.SourceContext.emitWatermark(Watermark mark)
Emits the given
Watermark . |
void |
ContinuousFileReaderOperator.processWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
Watermark |
BoundedOutOfOrdernessTimestampExtractor.getCurrentWatermark() |
Watermark |
AscendingTimestampExtractor.getCurrentWatermark()
Deprecated.
|
Modifier and Type | Method and Description |
---|---|
void |
InternalTimeServiceManagerImpl.advanceWatermark(Watermark watermark) |
void |
InternalTimeServiceManager.advanceWatermark(Watermark watermark)
Advances the Watermark of all managed
timer services ,
potentially firing event time timers. |
void |
TimestampedCollector.emitWatermark(Watermark mark) |
void |
Output.emitWatermark(Watermark mark)
Emits a
Watermark from an operator. |
void |
CountingOutput.emitWatermark(Watermark mark) |
void |
OnWatermarkCallback.onWatermark(KEY key,
Watermark watermark)
The action to be triggered upon reception of a watermark.
|
void |
ProcessOperator.processWatermark(Watermark mark) |
void |
StreamSink.processWatermark(Watermark mark) |
void |
AbstractInput.processWatermark(Watermark mark) |
void |
Input.processWatermark(Watermark mark)
Processes a
Watermark that arrived on the first input of this two-input operator. |
void |
AbstractStreamOperator.processWatermark(Watermark mark) |
void |
AbstractStreamOperatorV2.processWatermark(Watermark mark) |
void |
TwoInputStreamOperator.processWatermark1(Watermark mark)
Processes a
Watermark that arrived on the first input of this two-input operator. |
void |
AbstractStreamOperator.processWatermark1(Watermark mark) |
void |
TwoInputStreamOperator.processWatermark2(Watermark mark)
Processes a
Watermark that arrived on the second input of this two-input operator. |
void |
AbstractStreamOperator.processWatermark2(Watermark mark) |
protected void |
AbstractStreamOperatorV2.reportWatermark(Watermark mark,
int inputId) |
Modifier and Type | Method and Description |
---|---|
void |
AsyncWaitOperator.processWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
void |
CoProcessOperator.processWatermark(Watermark mark) |
void |
CoBroadcastWithNonKeyedOperator.processWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
void |
AbstractPythonFunctionOperator.processWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
void |
EmbeddedPythonCoProcessOperator.processWatermark(Watermark mark) |
void |
EmbeddedPythonProcessOperator.processWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
void |
ExternalPythonCoProcessOperator.processWatermark(Watermark mark) |
void |
ExternalPythonProcessOperator.processWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
void |
BatchExecutionInternalTimeServiceManager.advanceWatermark(Watermark watermark) |
Modifier and Type | Field and Description |
---|---|
static Watermark |
Watermark.MAX_WATERMARK
The watermark that signifies end-of-event-time.
|
static Watermark |
Watermark.UNINITIALIZED
The watermark that signifies is used before any actual watermark has been generated.
|
Modifier and Type | Method and Description |
---|---|
void |
PushingAsyncDataInput.DataOutput.emitWatermark(Watermark watermark) |
void |
RecordWriterOutput.emitWatermark(Watermark mark) |
void |
FinishedDataOutput.emitWatermark(Watermark watermark) |
Modifier and Type | Method and Description |
---|---|
void |
TimestampsAndWatermarksOperator.processWatermark(Watermark mark)
Override the base implementation to completely ignore watermarks propagated from upstream,
except for the "end of time" watermark.
|
Modifier and Type | Method and Description |
---|---|
Watermark |
StreamElement.asWatermark()
Casts this element into a Watermark.
|
Modifier and Type | Method and Description |
---|---|
void |
SourceOperatorStreamTask.AsyncDataOutputToOutput.emitWatermark(Watermark watermark) |
void |
FinishedOnRestoreMainOperatorOutput.emitWatermark(Watermark mark) |
void |
FinishedOnRestoreInput.processWatermark(Watermark watermark) |
Modifier and Type | Method and Description |
---|---|
void |
StatusWatermarkValve.inputWatermark(Watermark watermark,
int channelIndex,
PushingAsyncDataInput.DataOutput<?> output)
Feed a
Watermark into the valve. |
Modifier and Type | Method and Description |
---|---|
void |
TableStreamOperator.processWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
void |
LocalSlicingWindowAggOperator.processWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
void |
AbstractMapBundleOperator.processWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
void |
KeyedCoProcessOperatorWithWatermarkDelay.processWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
void |
OneInput.processWatermark(Watermark mark) |
void |
FirstInputOfTwoInput.processWatermark(Watermark mark) |
void |
SecondInputOfTwoInput.processWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
void |
OneInputStreamOperatorOutput.emitWatermark(Watermark mark) |
void |
CopyingSecondInputOfTwoInputStreamOperatorOutput.emitWatermark(Watermark mark) |
void |
BroadcastingOutput.emitWatermark(Watermark mark) |
void |
FirstInputOfTwoInputStreamOperatorOutput.emitWatermark(Watermark mark) |
void |
SecondInputOfTwoInputStreamOperatorOutput.emitWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
void |
SinkOperator.processWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
void |
InputConversionOperator.processWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
void |
SlicingWindowOperator.processWatermark(Watermark mark) |
Modifier and Type | Method and Description |
---|---|
Watermark |
PunctuatedWatermarkAssignerWrapper.checkAndGetNextWatermark(RowData row,
long extractedTimestamp) |
Watermark |
PeriodicWatermarkAssignerWrapper.getCurrentWatermark() |
Modifier and Type | Method and Description |
---|---|
void |
RowTimeMiniBatchAssginerOperator.processWatermark(Watermark mark) |
void |
WatermarkAssignerOperator.processWatermark(Watermark mark)
Override the base implementation to completely ignore watermarks propagated from upstream (we
rely only on the
WatermarkGenerator to emit watermarks from here). |
void |
ProcTimeMiniBatchAssignerOperator.processWatermark(Watermark mark)
Override the base implementation to completely ignore watermarks propagated from upstream (we
rely only on the
AssignerWithPeriodicWatermarks to emit watermarks from here). |
Modifier and Type | Method and Description |
---|---|
Watermark |
AscendingTimestamps.getWatermark() |
Watermark |
BoundedOutOfOrderTimestamps.getWatermark() |
abstract Watermark |
PeriodicWatermarkAssigner.getWatermark()
Returns the current watermark.
|
abstract Watermark |
PunctuatedWatermarkAssigner.getWatermark(Row row,
long timestamp)
Returns the watermark for the current row or null if no watermark should be generated.
|
Copyright © 2014–2024 The Apache Software Foundation. All rights reserved.