Modifier and Type | Method and Description |
---|---|
static <V> KafkaRecordDeserializationSchema<V> |
KafkaRecordDeserializationSchema.of(KafkaDeserializationSchema<V> kafkaDeserializationSchema)
Wraps a legacy
KafkaDeserializationSchema as the deserializer of the ConsumerRecords . |
Modifier and Type | Field and Description |
---|---|
protected KafkaDeserializationSchema<T> |
FlinkKafkaConsumerBase.deserializer
The schema to convert between Kafka's byte messages, and Flink's objects.
|
Constructor and Description |
---|
FlinkKafkaConsumer(List<String> topics,
KafkaDeserializationSchema<T> deserializer,
Properties props)
Deprecated.
Creates a new Kafka streaming source consumer.
|
FlinkKafkaConsumer(Pattern subscriptionPattern,
KafkaDeserializationSchema<T> deserializer,
Properties props)
Deprecated.
Creates a new Kafka streaming source consumer.
|
FlinkKafkaConsumer(String topic,
KafkaDeserializationSchema<T> deserializer,
Properties props)
Deprecated.
Creates a new Kafka streaming source consumer.
|
FlinkKafkaConsumerBase(List<String> topics,
Pattern topicPattern,
KafkaDeserializationSchema<T> deserializer,
long discoveryIntervalMillis,
boolean useMetrics)
Base constructor.
|
Modifier and Type | Class and Description |
---|---|
class |
KafkaDeserializationSchemaWrapper<T>
A simple wrapper for using the DeserializationSchema with the KafkaDeserializationSchema
interface.
|
Constructor and Description |
---|
KafkaFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
ProcessingTimeService processingTimeProvider,
long autoWatermarkInterval,
ClassLoader userCodeClassLoader,
String taskNameWithSubtasks,
KafkaDeserializationSchema<T> deserializer,
Properties kafkaProperties,
long pollTimeout,
MetricGroup subtaskMetricGroup,
MetricGroup consumerMetricGroup,
boolean useMetrics) |
KafkaShuffleFetcher(SourceFunction.SourceContext<T> sourceContext,
Map<KafkaTopicPartition,Long> assignedPartitionsWithInitialOffsets,
SerializedValue<WatermarkStrategy<T>> watermarkStrategy,
ProcessingTimeService processingTimeProvider,
long autoWatermarkInterval,
ClassLoader userCodeClassLoader,
String taskNameWithSubtasks,
KafkaDeserializationSchema<T> deserializer,
Properties kafkaProperties,
long pollTimeout,
MetricGroup subtaskMetricGroup,
MetricGroup consumerMetricGroup,
boolean useMetrics,
TypeSerializer<T> typeSerializer,
int producerParallelism) |
Modifier and Type | Interface and Description |
---|---|
interface |
KeyedDeserializationSchema<T>
Deprecated.
|
Modifier and Type | Class and Description |
---|---|
class |
JSONKeyValueDeserializationSchema
DeserializationSchema that deserializes a JSON String into an ObjectNode.
|
class |
TypeInformationKeyValueSerializationSchema<K,V>
A serialization and deserialization schema for Key Value Pairs that uses Flink's serialization
stack to transform typed from and to byte arrays.
|
Copyright © 2014–2023 The Apache Software Foundation. All rights reserved.