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 7d37d232e Use Shutdown Method in consumer runner
7d37d232e is described below
commit 7d37d232e15c9fffc986f33acb6ef3961a3d108d
Author: Mark Li <[email protected]>
AuthorDate: Wed Mar 1 03:06:32 2023 -0800
Use Shutdown Method in consumer runner
---
.../eventmesh/connector/kafka/consumer/ConsumerImpl.java | 10 ++--------
1 file changed, 2 insertions(+), 8 deletions(-)
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 92a61cf57..99c492d06 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
@@ -83,14 +83,8 @@ public class ConsumerImpl {
public synchronized void shutdown() {
if (this.started.compareAndSet(true, false)) {
// 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();
- }
- }
+ kafkaConsumerRunner.shutdown();
+ executorService.shutdown();
}
}
---------------------------------------------------------------------
To unsubscribe, e-mail: [email protected]
For additional commands, e-mail: [email protected]