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()) + "]"; }