Modifier and Type | Field and Description |
---|---|
protected DeserializationSchema<OUT> |
ConnectorSource.schema |
Constructor and Description |
---|
ConnectorSource(DeserializationSchema<OUT> schema) |
Constructor and Description |
---|
FlinkKafkaConsumer08(List<String> topics,
DeserializationSchema<T> deserializer,
Properties props)
Creates a new Kafka streaming source consumer for Kafka 0.8.x
This constructor allows passing multiple topics to the consumer.
|
FlinkKafkaConsumer08(String topic,
DeserializationSchema<T> valueDeserializer,
Properties props)
Creates a new Kafka streaming source consumer for Kafka 0.8.x
|
FlinkKafkaConsumer081(String topic,
DeserializationSchema<T> valueDeserializer,
Properties props)
Deprecated.
|
FlinkKafkaConsumer082(String topic,
DeserializationSchema<T> valueDeserializer,
Properties props)
Deprecated.
|
FlinkKafkaConsumer09(List<String> topics,
DeserializationSchema<T> deserializer,
Properties props)
Creates a new Kafka streaming source consumer for Kafka 0.9.x
This constructor allows passing multiple topics to the consumer.
|
FlinkKafkaConsumer09(String topic,
DeserializationSchema<T> valueDeserializer,
Properties props)
Creates a new Kafka streaming source consumer for Kafka 0.9.x
|
Modifier and Type | Field and Description |
---|---|
protected DeserializationSchema<OUT> |
RMQSource.schema |
Constructor and Description |
---|
RMQSource(String hostName,
Integer port,
String queueName,
boolean usesCorrelationId,
DeserializationSchema<OUT> deserializationSchema)
Creates a new RabbitMQ source.
|
RMQSource(String hostName,
Integer port,
String username,
String password,
String queueName,
boolean usesCorrelationId,
DeserializationSchema<OUT> deserializationSchema)
Creates a new RabbitMQ source.
|
RMQSource(String hostName,
String queueName,
boolean usesCorrelationId,
DeserializationSchema<OUT> deserializationSchema)
Creates a new RabbitMQ source.
|
RMQSource(String hostName,
String queueName,
DeserializationSchema<OUT> deserializationSchema)
Creates a new RabbitMQ source with at-least-once message processing guarantee when
checkpointing is enabled.
|
Modifier and Type | Class and Description |
---|---|
class |
SimpleStringSchema
Very simple serialization schema for strings.
|
class |
TypeInformationSerializationSchema<T>
A serialization and deserialization schema that uses Flink's serialization stack to
transform typed from and to byte arrays.
|
Constructor and Description |
---|
KeyedDeserializationSchemaWrapper(DeserializationSchema<T> deserializationSchema) |
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.