public class AppendBypassCompactWorkerOperator extends AppendCompactWorkerOperator<org.apache.flink.types.Either<Committable,UnawareAppendCompactionTask>>
AppendCompactWorkerOperator
to bypass Committable inputs.Modifier and Type | Class and Description |
---|---|
static class |
AppendBypassCompactWorkerOperator.Factory
StreamOperatorFactory of AppendBypassCompactWorkerOperator . |
unawareBucketCompactor
memoryPool
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 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–2025 The Apache Software Foundation. All rights reserved.