Package | Description |
---|---|
org.apache.flink.streaming.connectors.kafka | |
org.apache.flink.streaming.util.serialization |
Modifier and Type | Field and Description |
---|---|
protected KeyedSerializationSchema<IN> |
FlinkKafkaProducerBase.schema
(Serializable) SerializationSchema for turning objects used with Flink into
byte[] for Kafka.
|
Modifier and Type | Method and Description |
---|---|
static <T> FlinkKafkaProducer010.FlinkKafkaProducer010Configuration<T> |
FlinkKafkaProducer010.writeToKafkaWithTimestamps(DataStream<T> inStream,
String topicId,
KeyedSerializationSchema<T> serializationSchema,
Properties producerConfig)
Creates a FlinkKafkaProducer for a given topic.
|
static <T> FlinkKafkaProducer010.FlinkKafkaProducer010Configuration<T> |
FlinkKafkaProducer010.writeToKafkaWithTimestamps(DataStream<T> inStream,
String topicId,
KeyedSerializationSchema<T> serializationSchema,
Properties producerConfig,
KafkaPartitioner<T> customPartitioner)
Creates a FlinkKafkaProducer for a given topic.
|
Constructor and Description |
---|
FlinkKafkaProducer(String topicId,
KeyedSerializationSchema<IN> serializationSchema,
Properties producerConfig)
Deprecated.
|
FlinkKafkaProducer(String topicId,
KeyedSerializationSchema<IN> serializationSchema,
Properties producerConfig,
KafkaPartitioner customPartitioner)
Deprecated.
|
FlinkKafkaProducer(String brokerList,
String topicId,
KeyedSerializationSchema<IN> serializationSchema)
Deprecated.
|
FlinkKafkaProducer010(String topicId,
KeyedSerializationSchema<T> serializationSchema,
Properties producerConfig)
Creates a FlinkKafkaProducer for a given topic.
|
FlinkKafkaProducer010(String topicId,
KeyedSerializationSchema<T> serializationSchema,
Properties producerConfig,
KafkaPartitioner<T> customPartitioner)
Create Kafka producer
This constructor does not allow writing timestamps to Kafka, it follow approach (a) (see above)
|
FlinkKafkaProducer010(String brokerList,
String topicId,
KeyedSerializationSchema<T> serializationSchema)
Creates a FlinkKafkaProducer for a given topic.
|
FlinkKafkaProducer08(String topicId,
KeyedSerializationSchema<IN> serializationSchema,
Properties producerConfig)
Creates a FlinkKafkaProducer for a given topic.
|
FlinkKafkaProducer08(String topicId,
KeyedSerializationSchema<IN> serializationSchema,
Properties producerConfig,
KafkaPartitioner<IN> customPartitioner)
The main constructor for creating a FlinkKafkaProducer.
|
FlinkKafkaProducer08(String brokerList,
String topicId,
KeyedSerializationSchema<IN> serializationSchema)
Creates a FlinkKafkaProducer for a given topic.
|
FlinkKafkaProducer09(String topicId,
KeyedSerializationSchema<IN> serializationSchema,
Properties producerConfig)
Creates a FlinkKafkaProducer for a given topic.
|
FlinkKafkaProducer09(String topicId,
KeyedSerializationSchema<IN> serializationSchema,
Properties producerConfig,
KafkaPartitioner<IN> customPartitioner)
Creates a FlinkKafkaProducer for a given topic.
|
FlinkKafkaProducer09(String brokerList,
String topicId,
KeyedSerializationSchema<IN> serializationSchema)
Creates a FlinkKafkaProducer for a given topic.
|
FlinkKafkaProducerBase(String defaultTopicId,
KeyedSerializationSchema<IN> serializationSchema,
Properties producerConfig,
KafkaPartitioner<IN> customPartitioner)
The main constructor for creating a FlinkKafkaProducer.
|
Modifier and Type | Class and Description |
---|---|
class |
KeyedSerializationSchemaWrapper<T>
A simple wrapper for using the SerializationSchema 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.