Class | Description |
---|---|
AbstractFetcher<T,KPH> |
Base class for all fetchers, which implement the connections to Kafka brokers and
pull records from Kafka partitions.
|
ClosableBlockingQueue<E> |
A special form of blocking queue with two additions:
The queue can be closed atomically when empty.
|
ExceptionProxy |
A proxy that communicates exceptions between threads.
|
Kafka08Fetcher<T> |
A fetcher that fetches data from Kafka brokers via the Kafka 0.8 low-level consumer API.
|
KafkaTopicPartition |
Flink's description of a partition in a Kafka topic.
|
KafkaTopicPartitionLeader |
Serializable Topic Partition info with leader Node information.
|
KafkaTopicPartitionState<KPH> |
The state that the Flink Kafka Consumer holds for each Kafka partition.
|
KafkaTopicPartitionStateWithPeriodicWatermarks<T,KPH> |
A special version of the per-kafka-partition-state that additionally holds
a periodic watermark generator (and timestamp extractor) per partition.
|
KafkaTopicPartitionStateWithPunctuatedWatermarks<T,KPH> |
A special version of the per-kafka-partition-state that additionally holds
a periodic watermark generator (and timestamp extractor) per partition.
|
PeriodicOffsetCommitter |
A thread that periodically writes the current Kafka partition offsets to Zookeeper.
|
ZookeeperOffsetHandler |
Handler for committing Kafka offsets to Zookeeper and to retrieve them again.
|
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.