public class OneInputBroadcastWrapperOperator<IN,OUT> extends AbstractBroadcastWrapperOperator<OUT,org.apache.flink.streaming.api.operators.OneInputStreamOperator<IN,OUT>> implements org.apache.flink.streaming.api.operators.OneInputStreamOperator<IN,OUT>, org.apache.flink.streaming.api.operators.BoundedOneInput
OneInputStreamOperator
.containingTask, indexOfSubtask, metrics, numInputs, operatorFactory, output, parameters, stateHandler, streamConfig, timeServiceManager, wrappedOperator
Modifier and Type | Method and Description |
---|---|
void |
endInput() |
void |
processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<IN> streamRecord) |
void |
processLatencyMarker(org.apache.flink.streaming.runtime.streamrecord.LatencyMarker latencyMarker) |
void |
processWatermark(org.apache.flink.streaming.api.watermark.Watermark watermark) |
void |
processWatermarkStatus(org.apache.flink.streaming.runtime.watermarkstatus.WatermarkStatus watermarkStatus) |
void |
setKeyContextElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<IN> streamRecord) |
areBroadcastVariablesReady, close, endInputX, finish, getCurrentKey, getMetricGroup, getOperatorID, initializeState, initializeState, notifyCheckpointAborted, notifyCheckpointComplete, open, prepareSnapshotPreBarrier, processElementX, processWatermarkX, setCurrentKey, setKeyContextElement1, setKeyContextElement2, snapshotState, snapshotState
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
close, finish, getMetricGroup, getOperatorID, initializeState, open, prepareSnapshotPreBarrier, setKeyContextElement1, setKeyContextElement2, snapshotState
public void processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<IN> streamRecord) throws Exception
public void endInput() throws Exception
endInput
in interface org.apache.flink.streaming.api.operators.BoundedOneInput
Exception
public void processWatermark(org.apache.flink.streaming.api.watermark.Watermark watermark) throws Exception
public void processWatermarkStatus(org.apache.flink.streaming.runtime.watermarkstatus.WatermarkStatus watermarkStatus) throws Exception
public void processLatencyMarker(org.apache.flink.streaming.runtime.streamrecord.LatencyMarker latencyMarker) throws Exception
Copyright © 2019–2023 The Apache Software Foundation. All rights reserved.