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.