diff --git a/external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaSpout.java b/external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaSpout.java index 58347e334f5..ab34d4c5cd7 100644 --- a/external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaSpout.java +++ b/external/storm-kafka-client/src/main/java/org/apache/storm/kafka/spout/KafkaSpout.java @@ -162,6 +162,7 @@ private void initialize(Collection partitions) { KafkaSpoutMessageId msgId = msgIdIterator.next(); if (!partitionsSet.contains(msgId.getTopicPartition())) { msgIdIterator.remove(); + --numUncommittedOffsets; } }