Modifier and Type | Field and Description |
---|---|
protected KeyedDeserializationSchema<T> |
FlinkKafkaConsumerBase.deserializer
The schema to convert between Kafka#s byte messages, and Flink's objects
|
Constructor and Description |
---|
FlinkKafkaConsumer08(List<String> topics,
KeyedDeserializationSchema<T> deserializer,
Properties props)
Creates a new Kafka streaming source consumer for Kafka 0.8.x
This constructor allows passing multiple topics and a key/value deserialization schema.
|
FlinkKafkaConsumer08(String topic,
KeyedDeserializationSchema<T> deserializer,
Properties props)
Creates a new Kafka streaming source consumer for Kafka 0.8.x
This constructor allows passing a
KeyedDeserializationSchema for reading key/value
pairs, offsets, and topic names from Kafka. |
FlinkKafkaConsumer09(List<String> topics,
KeyedDeserializationSchema<T> deserializer,
Properties props)
Creates a new Kafka streaming source consumer for Kafka 0.9.x
This constructor allows passing multiple topics and a key/value deserialization schema.
|
FlinkKafkaConsumer09(String topic,
KeyedDeserializationSchema<T> deserializer,
Properties props)
Creates a new Kafka streaming source consumer for Kafka 0.9.x
This constructor allows passing a
KeyedDeserializationSchema for reading key/value
pairs, offsets, and topic names from Kafka. |
FlinkKafkaConsumerBase(KeyedDeserializationSchema<T> deserializer,
Properties props)
Creates a new Flink Kafka Consumer, using the given type of fetcher and offset handler.
|
Modifier and Type | Method and Description |
---|---|
<T> void |
LegacyFetcher.run(SourceFunction.SourceContext<T> sourceContext,
KeyedDeserializationSchema<T> deserializer,
HashMap<KafkaTopicPartition,Long> lastOffsets) |
<T> void |
Fetcher.run(SourceFunction.SourceContext<T> sourceContext,
KeyedDeserializationSchema<T> valueDeserializer,
HashMap<KafkaTopicPartition,Long> lastOffsets)
Starts fetch data from Kafka and emitting it into the stream.
|
Modifier and Type | Class and Description |
---|---|
class |
KeyedDeserializationSchemaWrapper<T>
A simple wrapper for using the DeserializationSchema with the KeyedDeserializationSchema
interface
|
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–2017 The Apache Software Foundation. All rights reserved.