public class KafkaActionUtils extends Object
Modifier and Type | Field and Description |
---|---|
static String |
PROPERTIES_PREFIX |
Constructor and Description |
---|
KafkaActionUtils() |
Modifier and Type | Method and Description |
---|---|
static org.apache.flink.connector.kafka.source.KafkaSource<CdcSourceRecord> |
buildKafkaSource(org.apache.flink.configuration.Configuration kafkaConfig,
org.apache.flink.streaming.connectors.kafka.KafkaDeserializationSchema<CdcSourceRecord> deserializationSchema) |
static String |
findOneTopic(org.apache.flink.configuration.Configuration kafkaConfig) |
static DataFormat |
getDataFormat(org.apache.flink.configuration.Configuration kafkaConfig) |
static MessageQueueSchemaUtils.ConsumerWrapper |
getKafkaEarliestConsumer(org.apache.flink.configuration.Configuration kafkaConfig,
org.apache.flink.streaming.connectors.kafka.KafkaDeserializationSchema<CdcSourceRecord> deserializationSchema) |
public static final String PROPERTIES_PREFIX
public static org.apache.flink.connector.kafka.source.KafkaSource<CdcSourceRecord> buildKafkaSource(org.apache.flink.configuration.Configuration kafkaConfig, org.apache.flink.streaming.connectors.kafka.KafkaDeserializationSchema<CdcSourceRecord> deserializationSchema)
public static DataFormat getDataFormat(org.apache.flink.configuration.Configuration kafkaConfig)
public static MessageQueueSchemaUtils.ConsumerWrapper getKafkaEarliestConsumer(org.apache.flink.configuration.Configuration kafkaConfig, org.apache.flink.streaming.connectors.kafka.KafkaDeserializationSchema<CdcSourceRecord> deserializationSchema)
public static String findOneTopic(org.apache.flink.configuration.Configuration kafkaConfig)
Copyright © 2023–2024 The Apache Software Foundation. All rights reserved.