From f4d3715317800597f449638958ce3bac66e498f1 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Thu, 12 Aug 2021 15:29:46 -0400 Subject: [PATCH] Recycle KafkaStreams Objects In the event Kafka Streams bindings are restarted (stop/start) using the actuator bindings endpoints, the underlying KafkaStreams objects are not recycled. After restarting, it still sees the previous KafkaStreams object. Addressing this issue. Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/1119 Resolves https://github.com/spring-cloud/spring-cloud-stream-binder-kafka/issues/1120 --- .../kafka/streams/GlobalKTableBinder.java | 13 ++++++++++- .../GlobalKTableBinderConfiguration.java | 5 +++-- .../binder/kafka/streams/KStreamBinder.java | 22 ++++++++++++++++++- .../streams/KStreamBinderConfiguration.java | 4 ++-- .../binder/kafka/streams/KTableBinder.java | 14 +++++++++++- .../streams/KTableBinderConfiguration.java | 5 +++-- .../kafka/streams/KafkaStreamsRegistry.java | 5 +++++ 7 files changed, 59 insertions(+), 9 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 5c22d4410..365ab9cd5 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 @@ -16,6 +16,7 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.kstream.GlobalKTable; import org.springframework.cloud.stream.binder.AbstractBinder; @@ -58,16 +59,18 @@ public class GlobalKTableBinder extends // @checkstyle:off private KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties = new KafkaStreamsExtendedBindingProperties(); + private final KafkaStreamsRegistry kafkaStreamsRegistry; // @checkstyle:on public GlobalKTableBinder( KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, KafkaTopicProvisioner kafkaTopicProvisioner, - KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue) { + KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, KafkaStreamsRegistry kafkaStreamsRegistry) { this.binderConfigurationProperties = binderConfigurationProperties; this.kafkaTopicProvisioner = kafkaTopicProvisioner; this.kafkaStreamsBindingInformationCatalogue = kafkaStreamsBindingInformationCatalogue; + this.kafkaStreamsRegistry = kafkaStreamsRegistry; } @Override @@ -97,9 +100,17 @@ public class GlobalKTableBinder extends return true; } + @Override + public synchronized void start() { + 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); } }; diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinderConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinderConfiguration.java index e18e338fc..6fd2f7a2d 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinderConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/GlobalKTableBinderConfiguration.java @@ -59,10 +59,11 @@ public class GlobalKTableBinderConfiguration { KafkaTopicProvisioner kafkaTopicProvisioner, KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, - @Qualifier("streamConfigGlobalProperties") Map streamConfigGlobalProperties) { + @Qualifier("streamConfigGlobalProperties") Map streamConfigGlobalProperties, + KafkaStreamsRegistry kafkaStreamsRegistry) { GlobalKTableBinder globalKTableBinder = new GlobalKTableBinder(binderConfigurationProperties, - kafkaTopicProvisioner, kafkaStreamsBindingInformationCatalogue); + kafkaTopicProvisioner, kafkaStreamsBindingInformationCatalogue, kafkaStreamsRegistry); globalKTableBinder.setKafkaStreamsExtendedBindingProperties( kafkaStreamsExtendedBindingProperties); return globalKTableBinder; 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 3ae3354cb..27f20a025 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 @@ -22,6 +22,7 @@ import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import org.apache.kafka.common.serialization.Serde; import org.apache.kafka.common.serialization.Serdes; +import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.kstream.KStream; import org.apache.kafka.streams.kstream.Produced; import org.apache.kafka.streams.processor.StreamPartitioner; @@ -78,16 +79,19 @@ class KStreamBinder extends private final KeyValueSerdeResolver keyValueSerdeResolver; + private final KafkaStreamsRegistry kafkaStreamsRegistry; + KStreamBinder(KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, KafkaTopicProvisioner kafkaTopicProvisioner, KafkaStreamsMessageConversionDelegate kafkaStreamsMessageConversionDelegate, KafkaStreamsBindingInformationCatalogue KafkaStreamsBindingInformationCatalogue, - KeyValueSerdeResolver keyValueSerdeResolver) { + KeyValueSerdeResolver keyValueSerdeResolver, KafkaStreamsRegistry kafkaStreamsRegistry) { this.binderConfigurationProperties = binderConfigurationProperties; this.kafkaTopicProvisioner = kafkaTopicProvisioner; this.kafkaStreamsMessageConversionDelegate = kafkaStreamsMessageConversionDelegate; this.kafkaStreamsBindingInformationCatalogue = KafkaStreamsBindingInformationCatalogue; this.keyValueSerdeResolver = keyValueSerdeResolver; + this.kafkaStreamsRegistry = kafkaStreamsRegistry; } @Override @@ -125,9 +129,17 @@ class KStreamBinder extends return true; } + @Override + public synchronized void start() { + 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); } }; @@ -178,9 +190,17 @@ class KStreamBinder extends return false; } + @Override + public synchronized void start() { + 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); } }; diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinderConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinderConfiguration.java index 26d094797..2341759c7 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinderConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KStreamBinderConfiguration.java @@ -59,10 +59,10 @@ public class KStreamBinderConfiguration { KafkaStreamsMessageConversionDelegate KafkaStreamsMessageConversionDelegate, KafkaStreamsBindingInformationCatalogue KafkaStreamsBindingInformationCatalogue, KeyValueSerdeResolver keyValueSerdeResolver, - KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties) { + KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, KafkaStreamsRegistry kafkaStreamsRegistry) { KStreamBinder kStreamBinder = new KStreamBinder(binderConfigurationProperties, kafkaTopicProvisioner, KafkaStreamsMessageConversionDelegate, - KafkaStreamsBindingInformationCatalogue, keyValueSerdeResolver); + KafkaStreamsBindingInformationCatalogue, keyValueSerdeResolver, kafkaStreamsRegistry); kStreamBinder.setKafkaStreamsExtendedBindingProperties( kafkaStreamsExtendedBindingProperties); return kStreamBinder; 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 acc1197d4..eb9cb849a 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 @@ -16,6 +16,7 @@ package org.springframework.cloud.stream.binder.kafka.streams; +import org.apache.kafka.streams.KafkaStreams; import org.apache.kafka.streams.kstream.KTable; import org.springframework.cloud.stream.binder.AbstractBinder; @@ -59,12 +60,15 @@ class KTableBinder extends // @checkstyle:on + private final KafkaStreamsRegistry kafkaStreamsRegistry; + KTableBinder(KafkaStreamsBinderConfigurationProperties binderConfigurationProperties, KafkaTopicProvisioner kafkaTopicProvisioner, - KafkaStreamsBindingInformationCatalogue KafkaStreamsBindingInformationCatalogue) { + KafkaStreamsBindingInformationCatalogue KafkaStreamsBindingInformationCatalogue, KafkaStreamsRegistry kafkaStreamsRegistry) { this.binderConfigurationProperties = binderConfigurationProperties; this.kafkaTopicProvisioner = kafkaTopicProvisioner; this.kafkaStreamsBindingInformationCatalogue = KafkaStreamsBindingInformationCatalogue; + this.kafkaStreamsRegistry = kafkaStreamsRegistry; } @Override @@ -97,9 +101,17 @@ class KTableBinder extends return true; } + @Override + public synchronized void start() { + 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); } }; diff --git a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinderConfiguration.java b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinderConfiguration.java index 50112f58f..ef0d5c94a 100644 --- a/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinderConfiguration.java +++ b/spring-cloud-stream-binder-kafka-streams/src/main/java/org/springframework/cloud/stream/binder/kafka/streams/KTableBinderConfiguration.java @@ -59,9 +59,10 @@ public class KTableBinderConfiguration { KafkaTopicProvisioner kafkaTopicProvisioner, KafkaStreamsExtendedBindingProperties kafkaStreamsExtendedBindingProperties, KafkaStreamsBindingInformationCatalogue kafkaStreamsBindingInformationCatalogue, - @Qualifier("streamConfigGlobalProperties") Map streamConfigGlobalProperties) { + @Qualifier("streamConfigGlobalProperties") Map streamConfigGlobalProperties, + KafkaStreamsRegistry kafkaStreamsRegistry) { KTableBinder kTableBinder = new KTableBinder(binderConfigurationProperties, - kafkaTopicProvisioner, kafkaStreamsBindingInformationCatalogue); + kafkaTopicProvisioner, kafkaStreamsBindingInformationCatalogue, kafkaStreamsRegistry); kTableBinder.setKafkaStreamsExtendedBindingProperties(kafkaStreamsExtendedBindingProperties); return kTableBinder; } 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 71808293d..e083dffc6 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 @@ -62,6 +62,11 @@ public class KafkaStreamsRegistry { this.streamsBuilderFactoryBeanMap.put(kafkaStreams, streamsBuilderFactoryBean); } + void unregisterKafkaStreams(KafkaStreams kafkaStreams) { + this.kafkaStreams.remove(kafkaStreams); + this.streamsBuilderFactoryBeanMap.remove(kafkaStreams); + } + /** * * @param kafkaStreams {@link KafkaStreams} object