| Package | Description |
|---|---|
| org.apache.paimon.flink.action.cdc | |
| org.apache.paimon.flink.action.cdc.kafka | |
| org.apache.paimon.flink.action.cdc.pulsar |
| Modifier and Type | Method and Description |
|---|---|
MessageQueueSchemaUtils.ConsumerWrapper |
SyncJobHandler.provideConsumer() |
| Modifier and Type | Method and Description |
|---|---|
static Schema |
MessageQueueSchemaUtils.getSchema(MessageQueueSchemaUtils.ConsumerWrapper consumer,
DataFormat dataFormat,
TypeMapping typeMapping)
Retrieves the Kafka schema for a given topic.
|
| Modifier and Type | Method and Description |
|---|---|
static MessageQueueSchemaUtils.ConsumerWrapper |
KafkaActionUtils.getKafkaEarliestConsumer(org.apache.flink.configuration.Configuration kafkaConfig,
org.apache.flink.streaming.connectors.kafka.KafkaDeserializationSchema<CdcSourceRecord> deserializationSchema) |
| Modifier and Type | Method and Description |
|---|---|
static MessageQueueSchemaUtils.ConsumerWrapper |
PulsarActionUtils.createPulsarConsumer(org.apache.flink.configuration.Configuration pulsarConfig,
org.apache.flink.api.common.serialization.DeserializationSchema<CdcSourceRecord> deserializationSchema)
Referenced to
PulsarPartitionSplitReader.createPulsarConsumer(org.apache.flink.connector.pulsar.source.enumerator.topic.TopicPartition). |
Copyright © 2023–2024 The Apache Software Foundation. All rights reserved.