Interface | Description |
---|---|
KafkaCommitCallback |
A callback interface that the source operator can implement to trigger custom actions when a
commit request completes, which should normally be triggered from checkpoint complete event.
|
Class | Description |
---|---|
AbstractFetcher<T,KPH> |
Base class for all fetchers, which implement the connections to Kafka brokers and pull records
from Kafka partitions.
|
AbstractPartitionDiscoverer |
Base class for all partition discoverers.
|
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.
|
KafkaDeserializationSchemaWrapper<T> |
A simple wrapper for using the DeserializationSchema with the KafkaDeserializationSchema
interface.
|
KafkaSerializationSchemaWrapper<T> |
An adapter from old style interfaces such as
SerializationSchema , FlinkKafkaPartitioner to the KafkaSerializationSchema . |
KafkaTopicPartition |
Flink's description of a partition in a Kafka topic.
|
KafkaTopicPartition.Comparator |
A
Comparator for KafkaTopicPartition s. |
KafkaTopicPartitionAssigner |
Utility for assigning Kafka partitions to consumer subtasks.
|
KafkaTopicPartitionLeader |
Serializable Topic Partition info with leader Node information.
|
KafkaTopicPartitionState<T,KPH> |
The state that the Flink Kafka Consumer holds for each Kafka partition.
|
KafkaTopicPartitionStateSentinel |
Magic values used to represent special offset states before partitions are actually read.
|
KafkaTopicPartitionStateWithWatermarkGenerator<T,KPH> |
A special version of the per-kafka-partition-state that additionally holds a
TimestampAssigner , WatermarkGenerator , an immediate WatermarkOutput , and a
deferred WatermarkOutput for this partition. |
KafkaTopicsDescriptor |
A Kafka Topics Descriptor describes how the consumer subscribes to Kafka topics - either a fixed
list of topics, or a topic pattern.
|
KeyedSerializationSchemaWrapper<T> |
A simple wrapper for using the SerializationSchema with the KeyedSerializationSchema interface.
|
SourceContextWatermarkOutputAdapter<T> |
A
WatermarkOutput that forwards calls to a SourceFunction.SourceContext . |
Exception | Description |
---|---|
AbstractPartitionDiscoverer.ClosedException |
Thrown if this discoverer was used to discover partitions after it was closed.
|
AbstractPartitionDiscoverer.WakeupException |
Signaling exception to indicate that an actual Kafka call was interrupted.
|
Copyright © 2014–2021 The Apache Software Foundation. All rights reserved.