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