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]