From 7d1aa4eadc440fbd054770af1574eabfef3edfad Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Mon, 18 Dec 2017 16:02:01 -0500 Subject: [PATCH] GH-522: Fix NPE in listener container (#524) Fixes https://github.com/spring-projects/spring-kafka/issues/522 * Polishing - PR Comment **Cherry-pick to 2.0.x & 1.3.x** --- .../KafkaMessageListenerContainer.java | 23 ++++++++++++------- 1 file changed, 15 insertions(+), 8 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 d40b2aff..2dbdd135 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 @@ -98,9 +98,9 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener private final TopicPartitionInitialOffset[] topicPartitions; - private ListenerConsumer listenerConsumer; + private volatile ListenerConsumer listenerConsumer; - private ListenableFuture listenerConsumerFuture; + private volatile ListenableFuture listenerConsumerFuture; private GenericMessageListener listener; @@ -180,11 +180,17 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener * either explicitly or by Kafka; may be null if not assigned yet. */ public Collection getAssignedPartitions() { - if (this.listenerConsumer.definedPartitions != null) { - return Collections.unmodifiableCollection(this.listenerConsumer.definedPartitions.keySet()); - } - else if (this.listenerConsumer.assignedPartitions != null) { - return Collections.unmodifiableCollection(this.listenerConsumer.assignedPartitions); + ListenerConsumer listenerConsumer = this.listenerConsumer; + if (listenerConsumer != null) { + if (listenerConsumer.definedPartitions != null) { + return Collections.unmodifiableCollection(listenerConsumer.definedPartitions.keySet()); + } + else if (listenerConsumer.assignedPartitions != null) { + return Collections.unmodifiableCollection(listenerConsumer.assignedPartitions); + } + else { + return null; + } } else { return null; @@ -294,7 +300,8 @@ public class KafkaMessageListenerContainer extends AbstractMessageListener public String toString() { return "KafkaMessageListenerContainer [id=" + getBeanName() + (this.clientIdSuffix != null ? ", clientIndex=" + this.clientIdSuffix : "") - + ", topicPartitions=" + getAssignedPartitions() + + ", topicPartitions=" + + (getAssignedPartitions() == null ? "none assigned" : getAssignedPartitions()) + "]"; }