[FLINK-4618] [kafka-connector] Minor improvements to comment and variable naming
This closes #2579 Project: http://git-wip-us.apache.org/repos/asf/flink/repo Commit: http://git-wip-us.apache.org/repos/asf/flink/commit/400c49ca Tree: http://git-wip-us.apache.org/repos/asf/flink/tree/400c49ca Diff: http://git-wip-us.apache.org/repos/asf/flink/diff/400c49ca Branch: refs/heads/release-1.1 Commit: 400c49ca02be088fbad77f018bea459e737269d3 Parents: bef0ebe Author: Tzu-Li (Gordon) Tai <[email protected]> Authored: Sat Oct 1 17:20:19 2016 +0800 Committer: Tzu-Li (Gordon) Tai <[email protected]> Committed: Sat Oct 1 18:16:51 2016 +0800 ---------------------------------------------------------------------- .../connectors/kafka/internal/Kafka09Fetcher.java | 14 ++++++-------- 1 file changed, 6 insertions(+), 8 deletions(-) ---------------------------------------------------------------------- http://git-wip-us.apache.org/repos/asf/flink/blob/400c49ca/flink-streaming-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/Kafka09Fetcher.java ---------------------------------------------------------------------- diff --git a/flink-streaming-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/Kafka09Fetcher.java b/flink-streaming-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/Kafka09Fetcher.java index 7e4177e..3c2cca3 100644 --- a/flink-streaming-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/Kafka09Fetcher.java +++ b/flink-streaming-connectors/flink-connector-kafka-0.9/src/main/java/org/apache/flink/streaming/connectors/kafka/internal/Kafka09Fetcher.java @@ -289,14 +289,12 @@ public class Kafka09Fetcher<T> extends AbstractFetcher<T, TopicPartition> implem Map<TopicPartition, OffsetAndMetadata> offsetsToCommit = new HashMap<>(partitions.length); for (KafkaTopicPartitionState<TopicPartition> partition : partitions) { - /* - * Increment offset by one, otherwise last record will be read again. This does not affect checkpoints/saved state. - * The offset is only read from Kafka/ZK on a fresh startup of a job, not restart or failure. See https://issues.apache.org/jira/browse/FLINK-4618 - */ - Long offset = offsets.get(partition.getKafkaTopicPartition()) + 1; - if (offset != null) { - offsetsToCommit.put(partition.getKafkaPartitionHandle(), new OffsetAndMetadata(offset)); - partition.setCommittedOffset(offset); + // committed offsets through the KafkaConsumer need to be 1 more than the last processed offset. + // This does not affect Flink's checkpoints/saved state. + Long offsetToCommit = offsets.get(partition.getKafkaTopicPartition()) + 1; + if (offsetToCommit != null) { + offsetsToCommit.put(partition.getKafkaPartitionHandle(), new OffsetAndMetadata(offsetToCommit)); + partition.setCommittedOffset(offsetToCommit); } }
