@Internal public static class RangeShuffle.LocalSampleOperator<T> extends org.apache.flink.table.runtime.operators.TableStreamOperator<org.apache.flink.api.java.tuple.Tuple3<Double,T,Integer>> implements org.apache.flink.streaming.api.operators.OneInputStreamOperator<org.apache.flink.api.java.tuple.Tuple2<T,Integer>,org.apache.flink.api.java.tuple.Tuple3<Double,T,Integer>>, org.apache.flink.streaming.api.operators.BoundedOneInput
See Sampler
.
Constructor and Description |
---|
LocalSampleOperator(int numSample) |
Modifier and Type | Method and Description |
---|---|
void |
endInput() |
void |
open() |
void |
processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<org.apache.flink.api.java.tuple.Tuple2<T,Integer>> streamRecord) |
computeMemorySize, processWatermark, useSplittableTimers
close, finish, getChainingStrategy, getContainingTask, getCurrentKey, getExecutionConfig, getInternalTimerService, getKeyedStateBackend, getKeyedStateStore, getMetricGroup, getOperatorConfig, getOperatorID, getOperatorName, getOperatorStateBackend, getOrCreateKeyedState, getPartitionedState, getPartitionedState, getProcessingTimeService, getRuntimeContext, getStateKeySelector1, getStateKeySelector2, getTimeServiceManager, getUserCodeClassloader, hasKeyContext1, hasKeyContext2, initializeState, initializeState, isUsingCustomRawKeyedState, notifyCheckpointAborted, notifyCheckpointComplete, prepareSnapshotPreBarrier, processLatencyMarker, processLatencyMarker1, processLatencyMarker2, processRecordAttributes, processRecordAttributes1, processRecordAttributes2, processWatermark1, processWatermark2, processWatermarkStatus, processWatermarkStatus1, processWatermarkStatus2, reportOrForwardLatencyMarker, setChainingStrategy, setCurrentKey, setKeyContextElement1, setKeyContextElement2, setMailboxExecutor, setProcessingTimeService, setup, snapshotState, snapshotState
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
setKeyContextElement
close, finish, getMetricGroup, getOperatorAttributes, getOperatorID, initializeState, prepareSnapshotPreBarrier, setKeyContextElement1, setKeyContextElement2, snapshotState
notifyCheckpointAborted, notifyCheckpointComplete
getCurrentKey, setCurrentKey
public void open() throws Exception
open
in interface org.apache.flink.streaming.api.operators.StreamOperator<org.apache.flink.api.java.tuple.Tuple3<Double,T,Integer>>
open
in class org.apache.flink.table.runtime.operators.TableStreamOperator<org.apache.flink.api.java.tuple.Tuple3<Double,T,Integer>>
Exception
public void processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<org.apache.flink.api.java.tuple.Tuple2<T,Integer>> streamRecord) throws Exception
public void endInput()
endInput
in interface org.apache.flink.streaming.api.operators.BoundedOneInput
Copyright © 2023–2024 The Apache Software Foundation. All rights reserved.