diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java index 3e5d82c866..885b5d311b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/processor/DefaultMessagingProcessor.java @@ -43,6 +43,7 @@ import org.apache.rocketmq.proxy.common.ProxyContext; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.config.ProxyConfig; import org.apache.rocketmq.proxy.service.ServiceManager; +import org.apache.rocketmq.proxy.service.ServiceManagerFactory; import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; import org.apache.rocketmq.proxy.service.route.ProxyTopicRouteData; import org.apache.rocketmq.proxy.service.transaction.TransactionId; @@ -94,7 +95,7 @@ public class DefaultMessagingProcessor extends AbstractStartAndShutdown implemen } public static DefaultMessagingProcessor createForLocalMode(BrokerController brokerController, RPCHook rpcHook) { - return new DefaultMessagingProcessor(ServiceManager.createForLocalMode(brokerController, rpcHook)); + return new DefaultMessagingProcessor(ServiceManagerFactory.createForLocalMode(brokerController, rpcHook)); } public static DefaultMessagingProcessor createForClusterMode() { @@ -102,7 +103,7 @@ public class DefaultMessagingProcessor extends AbstractStartAndShutdown implemen } public static DefaultMessagingProcessor createForClusterMode(RPCHook rpcHook) { - return new DefaultMessagingProcessor(ServiceManager.createForClusterMode(rpcHook)); + return new DefaultMessagingProcessor(ServiceManagerFactory.createForClusterMode(rpcHook)); } protected void init() { 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 61d8537a92..1fda36c9f3 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 @@ -29,6 +29,7 @@ import org.apache.rocketmq.broker.client.ProducerManager; import org.apache.rocketmq.common.constant.LoggerName; import org.apache.rocketmq.logging.InternalLogger; import org.apache.rocketmq.logging.InternalLoggerFactory; +import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.config.ProxyConfig; import org.apache.rocketmq.proxy.service.mqclient.MQClientAPIFactory; @@ -43,7 +44,7 @@ import org.apache.rocketmq.proxy.service.transaction.ClusterTransactionService; import org.apache.rocketmq.proxy.service.transaction.TransactionService; import org.apache.rocketmq.remoting.RPCHook; -public class ClusterServiceManager extends ServiceManager { +public class ClusterServiceManager extends AbstractStartAndShutdown implements ServiceManager { private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME); private final ClusterTransactionService clusterTransactionService; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/LocalServiceManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/LocalServiceManager.java index 1a57721feb..c2438c939c 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/LocalServiceManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/LocalServiceManager.java @@ -19,6 +19,7 @@ package org.apache.rocketmq.proxy.service; import org.apache.rocketmq.broker.BrokerController; import org.apache.rocketmq.broker.client.ConsumerManager; import org.apache.rocketmq.broker.client.ProducerManager; +import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown; import org.apache.rocketmq.proxy.service.message.LocalMessageService; import org.apache.rocketmq.proxy.service.message.MessageService; import org.apache.rocketmq.proxy.service.relay.LocalProxyRelayService; @@ -29,7 +30,7 @@ import org.apache.rocketmq.proxy.service.transaction.LocalTransactionService; import org.apache.rocketmq.proxy.service.transaction.TransactionService; import org.apache.rocketmq.remoting.RPCHook; -public class LocalServiceManager extends ServiceManager { +public class LocalServiceManager extends AbstractStartAndShutdown implements ServiceManager { private final BrokerController brokerController; private final TopicRouteService topicRouteService; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/ServiceManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/ServiceManager.java index 273a515dab..a51a1a1307 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/ServiceManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/ServiceManager.java @@ -16,43 +16,24 @@ */ package org.apache.rocketmq.proxy.service; -import org.apache.rocketmq.broker.BrokerController; import org.apache.rocketmq.broker.client.ConsumerManager; import org.apache.rocketmq.broker.client.ProducerManager; -import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown; +import org.apache.rocketmq.proxy.common.StartAndShutdown; import org.apache.rocketmq.proxy.service.message.MessageService; import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; import org.apache.rocketmq.proxy.service.route.TopicRouteService; import org.apache.rocketmq.proxy.service.transaction.TransactionService; -import org.apache.rocketmq.remoting.RPCHook; -public abstract class ServiceManager extends AbstractStartAndShutdown { +public interface ServiceManager extends StartAndShutdown { + MessageService getMessageService(); - public static ServiceManager createForLocalMode(BrokerController brokerController) { - return createForLocalMode(brokerController, null); - } + TopicRouteService getTopicRouteService(); - public static ServiceManager createForLocalMode(BrokerController brokerController, RPCHook rpcHook) { - return new LocalServiceManager(brokerController, rpcHook); - } + ProducerManager getProducerManager(); - public static ServiceManager createForClusterMode() { - return createForClusterMode(null); - } + ConsumerManager getConsumerManager(); - public static ServiceManager createForClusterMode(RPCHook rpcHook) { - return new ClusterServiceManager(rpcHook); - } + TransactionService getTransactionService(); - public abstract MessageService getMessageService(); - - public abstract TopicRouteService getTopicRouteService(); - - public abstract ProducerManager getProducerManager(); - - public abstract ConsumerManager getConsumerManager(); - - public abstract TransactionService getTransactionService(); - - public abstract ProxyRelayService getProxyOutService(); + ProxyRelayService getProxyOutService(); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/ServiceManagerFactory.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/ServiceManagerFactory.java new file mode 100644 index 0000000000..c186752788 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/ServiceManagerFactory.java @@ -0,0 +1,38 @@ +/* + * Licensed to the Apache Software Foundation (ASF) under one or more + * contributor license agreements. See the NOTICE file distributed with + * this work for additional information regarding copyright ownership. + * The ASF licenses this file to You under the Apache License, Version 2.0 + * (the "License"); you may not use this file except in compliance with + * the License. You may obtain a copy of the License at + * + * http://www.apache.org/licenses/LICENSE-2.0 + * + * Unless required by applicable law or agreed to in writing, software + * distributed under the License is distributed on an "AS IS" BASIS, + * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied. + * See the License for the specific language governing permissions and + * limitations under the License. + */ +package org.apache.rocketmq.proxy.service; + +import org.apache.rocketmq.broker.BrokerController; +import org.apache.rocketmq.remoting.RPCHook; + +public class ServiceManagerFactory { + public static ServiceManager createForLocalMode(BrokerController brokerController) { + return createForLocalMode(brokerController, null); + } + + public static ServiceManager createForLocalMode(BrokerController brokerController, RPCHook rpcHook) { + return new LocalServiceManager(brokerController, rpcHook); + } + + public static ServiceManager createForClusterMode() { + return createForClusterMode(null); + } + + public static ServiceManager createForClusterMode(RPCHook rpcHook) { + return new ClusterServiceManager(rpcHook); + } +}