From 9ccf9ce48b47af42d156830459733e33b10f390c Mon Sep 17 00:00:00 2001 From: Gary Russell Date: Sat, 12 Mar 2022 13:05:56 -0500 Subject: [PATCH] GH-2293: Cancel Subscription on Stop --- .../binder/reactorkafka/ReactorKafkaBinder.java | 12 ++++++++++++ 1 file changed, 12 insertions(+) diff --git a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinder.java b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinder.java index cec65422f..07be04529 100644 --- a/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinder.java +++ b/binders/kafka-binder/spring-cloud-stream-binder-kafka-reactive/src/main/java/org/springframework/cloud/stream/binder/reactorkafka/ReactorKafkaBinder.java @@ -30,6 +30,7 @@ import org.apache.kafka.clients.producer.ProducerRecord; import org.apache.kafka.clients.producer.RecordMetadata; import org.apache.kafka.common.serialization.ByteArrayDeserializer; import org.apache.kafka.common.serialization.ByteArraySerializer; +import org.reactivestreams.Subscription; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import reactor.core.publisher.Sinks; @@ -131,15 +132,26 @@ public class ReactorKafkaBinder private final KafkaReceiver receiver = KafkaReceiver.create(opts); + private volatile Subscription subscription; + @SuppressWarnings("unchecked") @Override protected void doStart() { Flux> flux = receiver .receive() + .doOnSubscribe(subs -> this.subscription = subs) .map(record -> (Message) converter.toMessage(record, null, null, null)); subscribeToPublisher(flux); } + @Override + protected synchronized void doStop() { + if (this.subscription != null) { + this.subscription.cancel(); + this.subscription = null; + } + } + }; }