public class AppendBypassCompactWorkerOperator extends AppendCompactWorkerOperator<org.apache.flink.types.Either<Committable,UnawareAppendCompactionTask>>
AppendCompactWorkerOperator
to bypass Committable inputs.unawareBucketCompactor
memoryPool
Constructor and Description |
---|
AppendBypassCompactWorkerOperator(FileStoreTable table,
String commitUser) |
Modifier and Type | Method and Description |
---|---|
void |
open() |
void |
processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<org.apache.flink.types.Either<Committable,UnawareAppendCompactionTask>> element) |
close, prepareCommit
endInput, prepareSnapshotPreBarrier, setup
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, processLatencyMarker, processLatencyMarker1, processLatencyMarker2, processRecordAttributes, processRecordAttributes1, processRecordAttributes2, processWatermark, processWatermark1, processWatermark2, processWatermarkStatus, processWatermarkStatus1, processWatermarkStatus2, reportOrForwardLatencyMarker, setChainingStrategy, setCurrentKey, setKeyContextElement1, setKeyContextElement2, setMailboxExecutor, setProcessingTimeService, snapshotState, snapshotState, useSplittableTimers
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
setKeyContextElement
finish, getMetricGroup, getOperatorAttributes, getOperatorID, initializeState, setKeyContextElement1, setKeyContextElement2, snapshotState
notifyCheckpointAborted, notifyCheckpointComplete
getCurrentKey, setCurrentKey
public AppendBypassCompactWorkerOperator(FileStoreTable table, String commitUser)
public void open() throws Exception
open
in interface org.apache.flink.streaming.api.operators.StreamOperator<Committable>
open
in class AppendCompactWorkerOperator<org.apache.flink.types.Either<Committable,UnawareAppendCompactionTask>>
Exception
public void processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<org.apache.flink.types.Either<Committable,UnawareAppendCompactionTask>> element) throws Exception
Exception
Copyright © 2023–2024 The Apache Software Foundation. All rights reserved.