public class OneInputAllRoundWrapperOperator<IN,OUT> extends AbstractAllRoundWrapperOperator<OUT,org.apache.flink.streaming.api.operators.OneInputStreamOperator<IN,OUT>> implements org.apache.flink.streaming.api.operators.OneInputStreamOperator<IterationRecord<IN>,IterationRecord<OUT>>, org.apache.flink.streaming.api.operators.BoundedOneInput
wrappedOperator
containingTask, epochWatermarkSupplier, epochWatermarkTracker, eventBroadcastOutput, iterationContext, metrics, operatorFactory, output, parameters, proxyOutput, streamConfig, uniqueSenderId
Constructor and Description |
---|
OneInputAllRoundWrapperOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<IterationRecord<OUT>> parameters,
org.apache.flink.streaming.api.operators.StreamOperatorFactory<OUT> operatorFactory) |
Modifier and Type | Method and Description |
---|---|
void |
endInput() |
void |
processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<IterationRecord<IN>> element) |
void |
processLatencyMarker(org.apache.flink.streaming.runtime.streamrecord.LatencyMarker latencyMarker) |
void |
processWatermark(org.apache.flink.streaming.api.watermark.Watermark mark) |
void |
processWatermarkStatus(org.apache.flink.streaming.runtime.watermarkstatus.WatermarkStatus watermarkStatus) |
void |
setKeyContextElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<IterationRecord<IN>> record) |
close, finish, getCurrentKey, getMetricGroup, getOperatorID, initializeState, notifyCheckpointAborted, notifyCheckpointComplete, onEpochWatermarkIncrement, open, prepareSnapshotPreBarrier, setCurrentKey, setKeyContextElement1, setKeyContextElement2, snapshotState
clearIterationContextRound, endInput, notifyEpochWatermarkIncrement, onEpochWatermarkEvent, setIterationContextRound
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
close, finish, getMetricGroup, getOperatorID, initializeState, open, prepareSnapshotPreBarrier, setKeyContextElement1, setKeyContextElement2, snapshotState
public OneInputAllRoundWrapperOperator(org.apache.flink.streaming.api.operators.StreamOperatorParameters<IterationRecord<OUT>> parameters, org.apache.flink.streaming.api.operators.StreamOperatorFactory<OUT> operatorFactory)
public void processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<IterationRecord<IN>> element) throws Exception
processElement
in interface org.apache.flink.streaming.api.operators.Input<IterationRecord<IN>>
Exception
public void processWatermark(org.apache.flink.streaming.api.watermark.Watermark mark) throws Exception
processWatermark
in interface org.apache.flink.streaming.api.operators.Input<IterationRecord<IN>>
Exception
public void processWatermarkStatus(org.apache.flink.streaming.runtime.watermarkstatus.WatermarkStatus watermarkStatus) throws Exception
processWatermarkStatus
in interface org.apache.flink.streaming.api.operators.Input<IterationRecord<IN>>
Exception
public void processLatencyMarker(org.apache.flink.streaming.runtime.streamrecord.LatencyMarker latencyMarker) throws Exception
processLatencyMarker
in interface org.apache.flink.streaming.api.operators.Input<IterationRecord<IN>>
Exception
public void setKeyContextElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<IterationRecord<IN>> record) throws Exception
setKeyContextElement
in interface org.apache.flink.streaming.api.operators.Input<IterationRecord<IN>>
setKeyContextElement
in interface org.apache.flink.streaming.api.operators.OneInputStreamOperator<IterationRecord<IN>,IterationRecord<OUT>>
Exception
Copyright © 2019–2023 The Apache Software Foundation. All rights reserved.