public class PulsarActionUtils extends Object
Modifier and Type | Field and Description |
---|---|
static org.apache.flink.configuration.ConfigOption<List<String>> |
TOPIC |
static org.apache.flink.configuration.ConfigOption<String> |
TOPIC_PATTERN |
static org.apache.flink.configuration.ConfigOption<String> |
VALUE_FORMAT |
Constructor and Description |
---|
PulsarActionUtils() |
Modifier and Type | Method and Description |
---|---|
static org.apache.flink.connector.pulsar.source.PulsarSource<CdcSourceRecord> |
buildPulsarSource(org.apache.flink.configuration.Configuration pulsarConfig,
org.apache.flink.api.common.serialization.DeserializationSchema<CdcSourceRecord> deserializationSchema) |
static MessageQueueSchemaUtils.ConsumerWrapper |
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) . |
static String |
findOneTopic(org.apache.flink.configuration.Configuration pulsarConfig) |
static DataFormat |
getDataFormat(org.apache.flink.configuration.Configuration pulsarConfig) |
public static final org.apache.flink.configuration.ConfigOption<String> VALUE_FORMAT
public static final org.apache.flink.configuration.ConfigOption<String> TOPIC_PATTERN
public static org.apache.flink.connector.pulsar.source.PulsarSource<CdcSourceRecord> buildPulsarSource(org.apache.flink.configuration.Configuration pulsarConfig, org.apache.flink.api.common.serialization.DeserializationSchema<CdcSourceRecord> deserializationSchema)
public static DataFormat getDataFormat(org.apache.flink.configuration.Configuration pulsarConfig)
public static MessageQueueSchemaUtils.ConsumerWrapper createPulsarConsumer(org.apache.flink.configuration.Configuration pulsarConfig, org.apache.flink.api.common.serialization.DeserializationSchema<CdcSourceRecord> deserializationSchema)
PulsarPartitionSplitReader.createPulsarConsumer(org.apache.flink.connector.pulsar.source.enumerator.topic.TopicPartition)
.public static String findOneTopic(org.apache.flink.configuration.Configuration pulsarConfig)
Copyright © 2023–2024 The Apache Software Foundation. All rights reserved.