public class InputOperator<T> extends org.apache.flink.streaming.api.operators.AbstractStreamOperator<IterationRecord<T>> implements org.apache.flink.streaming.api.operators.OneInputStreamOperator<T,IterationRecord<T>>
IterationRecord
.Constructor and Description |
---|
InputOperator() |
Modifier and Type | Method and Description |
---|---|
void |
open() |
void |
processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<T> streamRecord) |
void |
processWatermark(org.apache.flink.streaming.api.watermark.Watermark mark) |
close, finish, getChainingStrategy, getContainingTask, getCurrentKey, getExecutionConfig, getInternalTimerService, getKeyedStateBackend, getKeyedStateStore, getMetricGroup, getOperatorConfig, getOperatorID, getOperatorName, getOperatorStateBackend, getOrCreateKeyedState, getPartitionedState, getPartitionedState, getProcessingTimeService, getRuntimeContext, getTimeServiceManager, getUserCodeClassloader, initializeState, initializeState, isUsingCustomRawKeyedState, notifyCheckpointAborted, notifyCheckpointComplete, prepareSnapshotPreBarrier, processLatencyMarker, processLatencyMarker1, processLatencyMarker2, processWatermark1, processWatermark2, processWatermarkStatus, processWatermarkStatus1, processWatermarkStatus2, reportOrForwardLatencyMarker, setChainingStrategy, setCurrentKey, setKeyContextElement1, setKeyContextElement2, setProcessingTimeService, setup, snapshotState, snapshotState
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
setKeyContextElement
close, finish, getMetricGroup, getOperatorID, initializeState, prepareSnapshotPreBarrier, setKeyContextElement1, setKeyContextElement2, snapshotState
notifyCheckpointAborted, notifyCheckpointComplete
public void open() throws Exception
open
in interface org.apache.flink.streaming.api.operators.StreamOperator<IterationRecord<T>>
open
in class org.apache.flink.streaming.api.operators.AbstractStreamOperator<IterationRecord<T>>
Exception
public void processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<T> streamRecord) throws Exception
public void processWatermark(org.apache.flink.streaming.api.watermark.Watermark mark)
processWatermark
in interface org.apache.flink.streaming.api.operators.Input<T>
processWatermark
in class org.apache.flink.streaming.api.operators.AbstractStreamOperator<IterationRecord<T>>
Copyright © 2019–2023 The Apache Software Foundation. All rights reserved.