@Internal public class Kafka09PartitionDiscoverer extends AbstractPartitionDiscoverer
AbstractPartitionDiscoverer.ClosedException, AbstractPartitionDiscoverer.WakeupException
Constructor and Description |
---|
Kafka09PartitionDiscoverer(KafkaTopicsDescriptor topicsDescriptor,
int indexOfThisSubtask,
int numParallelSubtasks,
Properties kafkaProperties) |
Modifier and Type | Method and Description |
---|---|
protected void |
closeConnections()
Close all established connections.
|
protected List<KafkaTopicPartition> |
getAllPartitionsForTopics(List<String> topics)
Fetch the list of all partitions for a specific topics list from Kafka.
|
protected List<String> |
getAllTopics()
Fetch the list of all topics from Kafka.
|
protected void |
initializeConnections()
Establish the required connections in order to fetch topics and partitions metadata.
|
protected void |
wakeupConnections()
Attempt to eagerly wakeup from blocking calls to Kafka in
AbstractPartitionDiscoverer.getAllTopics()
and AbstractPartitionDiscoverer.getAllPartitionsForTopics(List) . |
close, discoverPartitions, open, setAndCheckDiscoveredPartition, wakeup
public Kafka09PartitionDiscoverer(KafkaTopicsDescriptor topicsDescriptor, int indexOfThisSubtask, int numParallelSubtasks, Properties kafkaProperties)
protected void initializeConnections()
AbstractPartitionDiscoverer
initializeConnections
in class AbstractPartitionDiscoverer
protected List<String> getAllTopics() throws AbstractPartitionDiscoverer.WakeupException
AbstractPartitionDiscoverer
getAllTopics
in class AbstractPartitionDiscoverer
AbstractPartitionDiscoverer.WakeupException
protected List<KafkaTopicPartition> getAllPartitionsForTopics(List<String> topics) throws AbstractPartitionDiscoverer.WakeupException
AbstractPartitionDiscoverer
getAllPartitionsForTopics
in class AbstractPartitionDiscoverer
AbstractPartitionDiscoverer.WakeupException
protected void wakeupConnections()
AbstractPartitionDiscoverer
AbstractPartitionDiscoverer.getAllTopics()
and AbstractPartitionDiscoverer.getAllPartitionsForTopics(List)
.
If the invocation indeed results in interrupting an actual blocking Kafka call, the implementations
of AbstractPartitionDiscoverer.getAllTopics()
and
AbstractPartitionDiscoverer.getAllPartitionsForTopics(List)
are responsible of throwing a
AbstractPartitionDiscoverer.WakeupException
.
wakeupConnections
in class AbstractPartitionDiscoverer
protected void closeConnections() throws Exception
AbstractPartitionDiscoverer
closeConnections
in class AbstractPartitionDiscoverer
Exception
Copyright © 2014–2019 The Apache Software Foundation. All rights reserved.