From 57a9d6398d887246a8f61330cc0a2ecda59a8596 Mon Sep 17 00:00:00 2001 From: cbb Date: Tue, 6 Jun 2017 22:15:18 +0800 Subject: [PATCH] Fix getHighestOffsetRecords * Add test case for "fix getHighestOffsetRecords" * Simple code style polishing **Cherry-pick to 1.2.x & 1.1.x** Conflicts: spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java spring-kafka/src/test/java/org/springframework/kafka/annotation/EnableKafkaIntegrationTests.java --- .../listener/KafkaMessageListenerContainer.java | 14 +++++--------- 1 file changed, 5 insertions(+), 9 deletions(-) diff --git a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java index 6542affb..ae2106f4 100644 --- a/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java +++ b/spring-kafka/src/main/java/org/springframework/kafka/listener/KafkaMessageListenerContainer.java @@ -982,18 +982,14 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener } private Collection> getHighestOffsetRecords(List> records) { - Map> highestOffsetMap = new HashMap<>(); - + Map> highestOffsetMap = new HashMap<>(); for (ConsumerRecord record : records) { - if (record != null) { - ConsumerRecord consumerRecord = highestOffsetMap.get(record.partition()); - - if (consumerRecord == null || record.offset() > consumerRecord.offset()) { - highestOffsetMap.put(record.partition(), record); - } + TopicPartition topicPartition = new TopicPartition(record.topic(), record.partition()); + ConsumerRecord consumerRecord = highestOffsetMap.get(topicPartition); + if (consumerRecord == null || record.offset() > consumerRecord.offset()) { + highestOffsetMap.put(topicPartition, record); } } - return highestOffsetMap.values(); }