From afe39bf78a59344c4f7bf1139344b98daf17a1be Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Tue, 24 Aug 2021 16:49:36 -0400 Subject: [PATCH] Kafka Streams binding lifecycle changes Start Kafka Streams bindgings only when they are not running. Similarly, stop them only if they are running. Without these guards in the bindings for KStream, KTable and GlobalKTable, it may cause NPE's due to the backing concurrent collections in KafkaStreamsRegistry not finding the proper KafkaStreams object, especially when the StreamsBuilderFactory bean is already stopped through the binder provided manager. --- .../kafka/streams/GlobalKTableBinder.java | 16 ++++++---- .../binder/kafka/streams/KStreamBinder.java | 32 ++++++++++++------- .../binder/kafka/streams/KTableBinder.java | 16 ++++++---- 3 files changed, 40 insertions(+), 24 deletions(-) diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinder.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinder.java index 365ab9cd5..418bda59a 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinder.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinder.java @@ -102,16 +102,20 @@ public class GlobalKTableBinder extends @Override public synchronized void start() { - super.start(); - GlobalKTableBinder.this.kafkaStreamsRegistry.registerKafkaStreams(streamsBuilderFactoryBean); + if (!streamsBuilderFactoryBean.isRunning()) { + super.start(); + GlobalKTableBinder.this.kafkaStreamsRegistry.registerKafkaStreams(streamsBuilderFactoryBean); + } } @Override public synchronized void stop() { - final KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams(); - super.stop(); - GlobalKTableBinder.this.kafkaStreamsRegistry.unregisterKafkaStreams(kafkaStreams); - KafkaStreamsBinderUtils.closeDlqProducerFactories(kafkaStreamsBindingInformationCatalogue, streamsBuilderFactoryBean); + if (streamsBuilderFactoryBean.isRunning()) { + final KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams(); + super.stop(); + GlobalKTableBinder.this.kafkaStreamsRegistry.unregisterKafkaStreams(kafkaStreams); + KafkaStreamsBinderUtils.closeDlqProducerFactories(kafkaStreamsBindingInformationCatalogue, streamsBuilderFactoryBean); + } } }; } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java index 27f20a025..e90cdbd1e 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinder.java @@ -131,16 +131,20 @@ class KStreamBinder extends @Override public synchronized void start() { - super.start(); - KStreamBinder.this.kafkaStreamsRegistry.registerKafkaStreams(streamsBuilderFactoryBean); + if (!streamsBuilderFactoryBean.isRunning()) { + super.start(); + KStreamBinder.this.kafkaStreamsRegistry.registerKafkaStreams(streamsBuilderFactoryBean); + } } @Override public synchronized void stop() { - final KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams(); - super.stop(); - KStreamBinder.this.kafkaStreamsRegistry.unregisterKafkaStreams(kafkaStreams); - KafkaStreamsBinderUtils.closeDlqProducerFactories(kafkaStreamsBindingInformationCatalogue, streamsBuilderFactoryBean); + if (streamsBuilderFactoryBean.isRunning()) { + final KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams(); + super.stop(); + KStreamBinder.this.kafkaStreamsRegistry.unregisterKafkaStreams(kafkaStreams); + KafkaStreamsBinderUtils.closeDlqProducerFactories(kafkaStreamsBindingInformationCatalogue, streamsBuilderFactoryBean); + } } }; } @@ -192,16 +196,20 @@ class KStreamBinder extends @Override public synchronized void start() { - super.start(); - KStreamBinder.this.kafkaStreamsRegistry.registerKafkaStreams(streamsBuilderFactoryBean); + if (!streamsBuilderFactoryBean.isRunning()) { + super.start(); + KStreamBinder.this.kafkaStreamsRegistry.registerKafkaStreams(streamsBuilderFactoryBean); + } } @Override public synchronized void stop() { - final KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams(); - super.stop(); - KStreamBinder.this.kafkaStreamsRegistry.unregisterKafkaStreams(kafkaStreams); - KafkaStreamsBinderUtils.closeDlqProducerFactories(kafkaStreamsBindingInformationCatalogue, streamsBuilderFactoryBean); + if (streamsBuilderFactoryBean.isRunning()) { + final KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams(); + super.stop(); + KStreamBinder.this.kafkaStreamsRegistry.unregisterKafkaStreams(kafkaStreams); + KafkaStreamsBinderUtils.closeDlqProducerFactories(kafkaStreamsBindingInformationCatalogue, streamsBuilderFactoryBean); + } } }; } diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java index eb9cb849a..c82a520b4 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinder.java @@ -103,16 +103,20 @@ class KTableBinder extends @Override public synchronized void start() { - super.start(); - KTableBinder.this.kafkaStreamsRegistry.registerKafkaStreams(streamsBuilderFactoryBean); + if (!streamsBuilderFactoryBean.isRunning()) { + super.start(); + KTableBinder.this.kafkaStreamsRegistry.registerKafkaStreams(streamsBuilderFactoryBean); + } } @Override public synchronized void stop() { - final KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams(); - super.stop(); - KTableBinder.this.kafkaStreamsRegistry.unregisterKafkaStreams(kafkaStreams); - KafkaStreamsBinderUtils.closeDlqProducerFactories(kafkaStreamsBindingInformationCatalogue, streamsBuilderFactoryBean); + if (streamsBuilderFactoryBean.isRunning()) { + final KafkaStreams kafkaStreams = streamsBuilderFactoryBean.getKafkaStreams(); + super.stop(); + KTableBinder.this.kafkaStreamsRegistry.unregisterKafkaStreams(kafkaStreams); + KafkaStreamsBinderUtils.closeDlqProducerFactories(kafkaStreamsBindingInformationCatalogue, streamsBuilderFactoryBean); + } } }; }