[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);
                        }
                }
 

Reply via email to