From b53df99c1b3ee9ba391ca602f66ff734f2325e58 Mon Sep 17 00:00:00 2001 From: Soby Chacko Date: Fri, 17 Feb 2023 12:11:04 -0500 Subject: [PATCH] Implement stop() in reader container --- .../reader/DefaultPulsarReaderListenerContainer.java | 10 +++++++++- 1 file changed, 9 insertions(+), 1 deletion(-) diff --git a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarReaderListenerContainer.java b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarReaderListenerContainer.java index 17dc6bc4..544bcee0 100644 --- a/spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarReaderListenerContainer.java +++ b/spring-pulsar/src/main/java/org/springframework/pulsar/reader/DefaultPulsarReaderListenerContainer.java @@ -16,6 +16,7 @@ package org.springframework.pulsar.reader; +import java.io.IOException; import java.util.concurrent.CompletableFuture; import java.util.concurrent.CountDownLatch; import java.util.concurrent.TimeUnit; @@ -98,7 +99,14 @@ public class DefaultPulsarReaderListenerContainer extends AbstractPulsarReade @Override protected void doStop() { - + setRunning(false); + try { + this.logger.info("Closing this consumer."); + this.internalAsyncReader.get().reader.close(); + } + catch (IOException e) { + this.logger.error(e, () -> "Error closing Pulsar Client."); + } } private void publishReaderStartingEvent() {