Package | Description |
---|---|
org.apache.flink.streaming.connectors.kafka | |
org.apache.flink.streaming.connectors.kafka.internals |
Modifier and Type | Field and Description |
---|---|
protected HashMap<KafkaTopicPartition,Long> |
FlinkKafkaConsumerBase.offsetsState
The offsets of the last returned elements
|
protected HashMap<KafkaTopicPartition,Long> |
FlinkKafkaConsumerBase.restoreToOffset
The offsets to restore to, if the consumer restores state from a checkpoint
|
Modifier and Type | Method and Description |
---|---|
static List<KafkaTopicPartition> |
FlinkKafkaConsumer09.convertToFlinkKafkaTopicPartition(List<org.apache.kafka.common.PartitionInfo> partitions)
Converts a list of Kafka PartitionInfo's to Flink's KafkaTopicPartition (which are serializable)
|
HashMap<KafkaTopicPartition,Long> |
FlinkKafkaConsumerBase.snapshotState(long checkpointId,
long checkpointTimestamp) |
Modifier and Type | Method and Description |
---|---|
protected void |
FlinkKafkaConsumer09.commitOffsets(HashMap<KafkaTopicPartition,Long> checkpointOffsets) |
protected abstract void |
FlinkKafkaConsumerBase.commitOffsets(HashMap<KafkaTopicPartition,Long> checkpointOffsets) |
protected void |
FlinkKafkaConsumer08.commitOffsets(HashMap<KafkaTopicPartition,Long> toCommit)
Utility method to commit offsets.
|
static Map<org.apache.kafka.common.TopicPartition,org.apache.kafka.clients.consumer.OffsetAndMetadata> |
FlinkKafkaConsumer09.convertToCommitMap(HashMap<KafkaTopicPartition,Long> checkpointOffsets) |
static List<org.apache.kafka.common.TopicPartition> |
FlinkKafkaConsumer09.convertToKafkaTopicPartition(List<KafkaTopicPartition> partitions) |
static void |
FlinkKafkaConsumerBase.logPartitionInfo(List<KafkaTopicPartition> partitionInfos)
Method to log partition information.
|
void |
FlinkKafkaConsumerBase.restoreState(HashMap<KafkaTopicPartition,Long> restoredOffsets) |
Modifier and Type | Method and Description |
---|---|
KafkaTopicPartition |
KafkaTopicPartitionLeader.getTopicPartition() |
Modifier and Type | Method and Description |
---|---|
static List<KafkaTopicPartition> |
KafkaTopicPartition.dropLeaderData(List<KafkaTopicPartitionLeader> partitionInfos) |
Map<KafkaTopicPartition,Long> |
ZookeeperOffsetHandler.getOffsets(List<KafkaTopicPartition> partitions) |
Map<KafkaTopicPartition,Long> |
OffsetHandler.getOffsets(List<KafkaTopicPartition> partitions)
Positions the given fetcher to the initial read offsets where the stream consumption
will start from.
|
Modifier and Type | Method and Description |
---|---|
void |
ZookeeperOffsetHandler.commit(Map<KafkaTopicPartition,Long> offsetsToCommit) |
void |
OffsetHandler.commit(Map<KafkaTopicPartition,Long> offsetsToCommit)
Commits the given offset for the partitions.
|
Map<KafkaTopicPartition,Long> |
ZookeeperOffsetHandler.getOffsets(List<KafkaTopicPartition> partitions) |
Map<KafkaTopicPartition,Long> |
OffsetHandler.getOffsets(List<KafkaTopicPartition> partitions)
Positions the given fetcher to the initial read offsets where the stream consumption
will start from.
|
<T> void |
LegacyFetcher.run(SourceFunction.SourceContext<T> sourceContext,
KeyedDeserializationSchema<T> deserializer,
HashMap<KafkaTopicPartition,Long> lastOffsets) |
<T> void |
Fetcher.run(SourceFunction.SourceContext<T> sourceContext,
KeyedDeserializationSchema<T> valueDeserializer,
HashMap<KafkaTopicPartition,Long> lastOffsets)
Starts fetch data from Kafka and emitting it into the stream.
|
static String |
KafkaTopicPartition.toString(List<KafkaTopicPartition> partitions) |
static String |
KafkaTopicPartition.toString(Map<KafkaTopicPartition,Long> map) |
Constructor and Description |
---|
KafkaTopicPartitionLeader(KafkaTopicPartition topicPartition,
org.apache.kafka.common.Node leader) |
Constructor and Description |
---|
LegacyFetcher(Map<KafkaTopicPartition,Long> initialPartitionsToRead,
Properties props,
String taskName,
ClassLoader userCodeClassloader)
Create a LegacyFetcher instance.
|
Copyright © 2014–2017 The Apache Software Foundation. All rights reserved.