This is an automated email from the ASF dual-hosted git repository.

markli pushed a commit to branch Kafka-Shutdown
in repository https://gitbox.apache.org/repos/asf/incubator-eventmesh.git


The following commit(s) were added to refs/heads/Kafka-Shutdown by this push:
     new 0f517b62a Consumer Proper Shutdown
0f517b62a is described below

commit 0f517b62ab5243b2c3299f41dc44f4b6157b8adc
Author: Mark Li <[email protected]>
AuthorDate: Wed Mar 1 02:38:43 2023 -0800

    Consumer Proper Shutdown
---
 .../eventmesh/connector/kafka/consumer/ConsumerImpl.java       | 10 +++++++++-
 1 file changed, 9 insertions(+), 1 deletion(-)

diff --git 
a/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/consumer/ConsumerImpl.java
 
b/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/consumer/ConsumerImpl.java
index eb370c459..92a61cf57 100644
--- 
a/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/consumer/ConsumerImpl.java
+++ 
b/eventmesh-connector-plugin/eventmesh-connector-kafka/src/main/java/org/apache/eventmesh/connector/kafka/consumer/ConsumerImpl.java
@@ -82,7 +82,15 @@ public class ConsumerImpl {
 
     public synchronized void shutdown() {
         if (this.started.compareAndSet(true, false)) {
-            this.kafkaConsumer.close();
+            // Shutdown the executor and interrupt any running tasks
+            List<Runnable> tasks = executorService.shutdownNow();
+
+            // Call a shutdown on each thread
+            for (Runnable task : tasks) {
+                if (task instanceof KafkaConsumerRunner) {
+                    ((KafkaConsumerRunner) task).shutdown();
+                }
+            }
         }
     }
 


---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]

Reply via email to