Package | Description |
---|---|
org.apache.paimon.flink.sink.cdc |
Modifier and Type | Method and Description |
---|---|
static CdcRecord |
CdcRecord.emptyRecord() |
CdcRecord |
CdcRecord.fieldNameLowerCase() |
static CdcRecord |
CdcRecordUtils.fromGenericRow(GenericRow row,
List<String> fieldNames) |
CdcRecord |
CdcMultiplexRecord.record() |
CdcRecord |
RichCdcRecord.toCdcRecord() |
Modifier and Type | Method and Description |
---|---|
static org.apache.flink.streaming.api.datastream.DataStream<CdcRecord> |
CaseSensitiveUtils.cdcRecordConvert(Catalog.Loader catalogLoader,
org.apache.flink.streaming.api.datastream.DataStream<CdcRecord> input) |
protected KeyAndBucketExtractor<CdcRecord> |
CdcRecordChannelComputer.createExtractor() |
protected KeyAndBucketExtractor<CdcRecord> |
CdcDynamicBucketSink.createExtractor(TableSchema schema) |
static org.apache.flink.util.OutputTag<CdcRecord> |
CdcMultiTableParsingProcessFunction.createRecordOutputTag(String tableName) |
protected org.apache.flink.streaming.api.operators.OneInputStreamOperator<CdcRecord,Committable> |
CdcFixedBucketSink.createWriteOperator(StoreSinkWrite.Provider writeProvider,
String commitUser) |
protected org.apache.flink.streaming.api.operators.OneInputStreamOperator<CdcRecord,Committable> |
CdcUnawareBucketSink.createWriteOperator(StoreSinkWrite.Provider writeProvider,
String commitUser) |
protected org.apache.flink.streaming.api.operators.OneInputStreamOperator<org.apache.flink.api.java.tuple.Tuple2<CdcRecord,Integer>,Committable> |
CdcDynamicBucketSink.createWriteOperator(StoreSinkWrite.Provider writeProvider,
String commitUser) |
List<CdcRecord> |
RichEventParser.parseRecords() |
List<CdcRecord> |
EventParser.parseRecords()
Parse records from event.
|
List<CdcRecord> |
RichCdcMultiplexRecordEventParser.parseRecords() |
Modifier and Type | Method and Description |
---|---|
static CdcMultiplexRecord |
CdcMultiplexRecord.fromCdcRecord(String databaseName,
String tableName,
CdcRecord record) |
BinaryRow |
CdcRecordPartitionKeyExtractor.partition(CdcRecord record) |
static GenericRow |
CdcRecordUtils.projectAsInsert(CdcRecord record,
List<DataField> dataFields)
Project
fields to a GenericRow . |
void |
CdcRecordKeyAndBucketExtractor.setRecord(CdcRecord record) |
static Optional<GenericRow> |
CdcRecordUtils.toGenericRow(CdcRecord record,
List<DataField> dataFields)
Convert
fields to a GenericRow . |
BinaryRow |
CdcRecordPartitionKeyExtractor.trimmedPrimaryKey(CdcRecord record) |
Modifier and Type | Method and Description |
---|---|
static org.apache.flink.streaming.api.datastream.DataStream<CdcRecord> |
CaseSensitiveUtils.cdcRecordConvert(Catalog.Loader catalogLoader,
org.apache.flink.streaming.api.datastream.DataStream<CdcRecord> input) |
void |
CdcRecordStoreWriteOperator.processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<CdcRecord> element) |
void |
CdcUnawareBucketWriteOperator.processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<CdcRecord> element) |
void |
CdcDynamicBucketWriteOperator.processElement(org.apache.flink.streaming.runtime.streamrecord.StreamRecord<org.apache.flink.api.java.tuple.Tuple2<CdcRecord,Integer>> element) |
void |
CdcParsingProcessFunction.processElement(T raw,
org.apache.flink.streaming.api.functions.ProcessFunction.Context context,
org.apache.flink.util.Collector<CdcRecord> collector) |
Constructor and Description |
---|
CdcMultiplexRecord(String databaseName,
String tableName,
CdcRecord record) |
RichCdcMultiplexRecord(String databaseName,
String tableName,
List<DataField> fields,
List<String> primaryKeys,
CdcRecord cdcRecord) |
RichCdcRecord(CdcRecord cdcRecord,
List<DataField> fields) |
Copyright © 2023–2024 The Apache Software Foundation. All rights reserved.