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

mikexue pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/eventmesh.git


The following commit(s) were added to refs/heads/master by this push:
     new e7081caf6 [ISSUE #3707]fix thread leak with rocketmq consumer start 
and shutdown frequently
e7081caf6 is described below

commit e7081caf6e4d47d508c41f1d3756b56cd761e81e
Author: willimpo <[email protected]>
AuthorDate: Mon Apr 17 16:47:31 2023 +0800

    [ISSUE #3707]fix thread leak with rocketmq consumer start and shutdown 
frequently
---
 .../client/impl/consumer/ConsumeMessageConcurrentlyService.java       | 4 +++-
 1 file changed, 3 insertions(+), 1 deletion(-)

diff --git 
a/eventmesh-storage-plugin/eventmesh-storage-rocketmq/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyService.java
 
b/eventmesh-storage-plugin/eventmesh-storage-rocketmq/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyService.java
index 78b1257ef..80fd263b7 100644
--- 
a/eventmesh-storage-plugin/eventmesh-storage-rocketmq/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyService.java
+++ 
b/eventmesh-storage-plugin/eventmesh-storage-rocketmq/src/main/java/org/apache/rocketmq/client/impl/consumer/ConsumeMessageConcurrentlyService.java
@@ -106,7 +106,9 @@ public class ConsumeMessageConcurrentlyService implements 
ConsumeMessageService
 
     @Override
     public void shutdown(long awaitTerminateMillis) {
-
+        this.scheduledExecutorService.shutdown();
+        
org.apache.rocketmq.common.utils.ThreadUtils.shutdownGracefully(this.consumeExecutor,
 awaitTerminateMillis, TimeUnit.MILLISECONDS);
+        this.cleanExpireMsgExecutors.shutdown();
     }
 
     public void shutdown() {


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

Reply via email to