public class CdcDynamicBucketSink extends CdcDynamicBucketSinkBase<CdcRecord>
Constructor and Description |
---|
CdcDynamicBucketSink(FileStoreTable table) |
Modifier and Type | Method and Description |
---|---|
protected KeyAndBucketExtractor<CdcRecord> |
createExtractor(TableSchema schema) |
protected org.apache.flink.streaming.api.operators.OneInputStreamOperator<org.apache.flink.api.java.tuple.Tuple2<CdcRecord,Integer>,Committable> |
createWriteOperator(StoreSinkWrite.Provider writeProvider,
String commitUser) |
assignerChannelComputer, channelComputer2, extractorFunction
build, createHashBucketAssignerOperator
createCommittableStateManager, createCommitterFactory
assertBatchAdaptiveParallelism, assertBatchAdaptiveParallelism, assertStreamingConfiguration, configureGlobalCommitter, doCommit, doWrite, isStreaming, isStreaming, sinkFrom, sinkFrom
public CdcDynamicBucketSink(FileStoreTable table)
protected KeyAndBucketExtractor<CdcRecord> createExtractor(TableSchema schema)
createExtractor
in class CdcDynamicBucketSinkBase<CdcRecord>
protected org.apache.flink.streaming.api.operators.OneInputStreamOperator<org.apache.flink.api.java.tuple.Tuple2<CdcRecord,Integer>,Committable> createWriteOperator(StoreSinkWrite.Provider writeProvider, String commitUser)
createWriteOperator
in class FlinkSink<org.apache.flink.api.java.tuple.Tuple2<CdcRecord,Integer>>
Copyright © 2023–2024 The Apache Software Foundation. All rights reserved.