public class CdcUnawareBucketWriteOperator extends CdcRecordStoreWriteOperator
PrepareCommitOperator
to write CdcRecord
to unaware-bucket mode table.Modifier and Type | Class and Description |
---|---|
static class |
CdcUnawareBucketWriteOperator.Factory
StreamOperatorFactory of CdcUnawareBucketWriteOperator . |
MAX_RETRY_NUM_TIMES, RETRY_SLEEP_TIME, SKIP_CORRUPT_RECORD
table, write
memoryPool
Modifier and Type | Method and Description |
---|---|
void |
processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<CdcRecord> element) |
containLogSystem, initializeState
close, createState, getCommitUser, getWrite, prepareCommit, processRecordAttributes, 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, 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 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–2025 The Apache Software Foundation. All rights reserved.