diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java index b4fbe48b4..27c1cc0a0 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsBinderSupportAutoConfiguration.java @@ -61,6 +61,7 @@ import org.springframework.core.env.MapPropertySource; import org.springframework.kafka.config.KafkaStreamsConfiguration; import org.springframework.kafka.core.CleanupConfig; import org.springframework.kafka.streams.RecoveringDeserializationExceptionHandler; +import org.springframework.lang.Nullable; import org.springframework.messaging.converter.CompositeMessageConverter; import org.springframework.util.ObjectUtils; import org.springframework.util.StringUtils; @@ -344,7 +345,7 @@ public class KafkaStreamsBinderSupportAutoConfiguration { } @Bean - public KafkaStreamsRegistry kafkaStreamsRegistry(KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics) { + public KafkaStreamsRegistry kafkaStreamsRegistry(@Nullable KafkaStreamsBinderMetrics kafkaStreamsBinderMetrics) { return new KafkaStreamsRegistry(kafkaStreamsBinderMetrics); } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsRegistry.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsRegistry.java index 88f4c28ab..79f65986d 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsRegistry.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KafkaStreamsRegistry.java @@ -52,7 +52,9 @@ class KafkaStreamsRegistry { * @param kafkaStreams {@link KafkaStreams} object created in the application */ void registerKafkaStreams(KafkaStreams kafkaStreams) { - this.kafkaStreamsBinderMetrics.addMetrics(kafkaStreams); + if (this.kafkaStreamsBinderMetrics != null) { + this.kafkaStreamsBinderMetrics.addMetrics(kafkaStreams); + } this.kafkaStreams.add(kafkaStreams); }