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