From bb5e55df163d34c8e19d6f08910a7d055f86e559 Mon Sep 17 00:00:00 2001 From: "kaiyi.lk" Date: Wed, 18 May 2022 17:36:58 +0800 Subject: [PATCH] [ISSUE #3949] add shutdown in ClusterMetadataService --- .../rocketmq/proxy/service/ClusterServiceManager.java | 3 ++- .../proxy/service/metadata/ClusterMetadataService.java | 6 ++++++ 2 files changed, 8 insertions(+), 1 deletion(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/ClusterServiceManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/ClusterServiceManager.java index 74b0052b03..91604cb13b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/ClusterServiceManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/ClusterServiceManager.java @@ -55,7 +55,7 @@ public class ClusterServiceManager extends AbstractStartAndShutdown implements S private final TopicRouteService topicRouteService; private final MessageService messageService; private final ProxyRelayService proxyRelayService; - private final MetadataService metadataService; + private final ClusterMetadataService metadataService; private final ScheduledExecutorService scheduledExecutorService; private final MQClientAPIFactory messagingClientAPIFactory; @@ -107,6 +107,7 @@ public class ClusterServiceManager extends AbstractStartAndShutdown implements S this.appendStartAndShutdown(this.operationClientAPIFactory); this.appendStartAndShutdown(this.topicRouteService); this.appendStartAndShutdown(this.clusterTransactionService); + this.appendStartAndShutdown(this.metadataService); } @Override diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/ClusterMetadataService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/ClusterMetadataService.java index 3be7b051b9..a5b18636d6 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/ClusterMetadataService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/ClusterMetadataService.java @@ -74,6 +74,12 @@ public class ClusterMetadataService extends AbstractStartAndShutdown implements .maximumSize(config.getSubscriptionGroupConfigCacheMaxNum()) .refreshAfterWrite(config.getSubscriptionGroupConfigCacheExpiredInSeconds(), TimeUnit.SECONDS) .build(new ClusterSubscriptionGroupConfigCacheLoader()); + + this.init(); + } + + protected void init() { + this.appendShutdown(this.cacheRefreshExecutor::shutdown); } @Override