public class CdcUnawareBucketWriteOperator extends CdcRecordStoreWriteOperator
PrepareCommitOperator
to write CdcRecord
to unaware-bucket mode table.RETRY_SLEEP_TIME
table, write
memoryPool
Constructor and Description |
---|
CdcUnawareBucketWriteOperator(FileStoreTable table,
StoreSinkWrite.Provider storeSinkWriteProvider,
String initialCommitUser) |
Modifier and Type | Method and Description |
---|---|
void |
processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<CdcRecord> element) |
containLogSystem, initializeState
close, getWrite, prepareCommit, snapshotState
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, isUsingCustomRawKeyedState, notifyCheckpointAborted, notifyCheckpointComplete, open, processLatencyMarker, processLatencyMarker1, processLatencyMarker2, processRecordAttributes, processRecordAttributes1, processRecordAttributes2, processWatermark, processWatermark1, processWatermark2, processWatermarkStatus, processWatermarkStatus1, processWatermarkStatus2, reportOrForwardLatencyMarker, setChainingStrategy, setCurrentKey, setKeyContextElement1, setKeyContextElement2, setMailboxExecutor, setProcessingTimeService, snapshotState, useSplittableTimers
clone, equals, finalize, getClass, hashCode, notify, notifyAll, toString, wait, wait, wait
setKeyContextElement
finish, getMetricGroup, getOperatorAttributes, getOperatorID, initializeState, open, setKeyContextElement1, setKeyContextElement2, snapshotState
notifyCheckpointAborted, notifyCheckpointComplete
getCurrentKey, setCurrentKey
public CdcUnawareBucketWriteOperator(FileStoreTable table, StoreSinkWrite.Provider storeSinkWriteProvider, String initialCommitUser)
public void processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<CdcRecord> element) throws Exception
processElement
in interface org.apache.flink.streaming.api.operators.Input<CdcRecord>
processElement
in class CdcRecordStoreWriteOperator
Exception
Copyright © 2023–2024 The Apache Software Foundation. All rights reserved.