diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/TopicConfigCache.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/TopicConfigCache.java deleted file mode 100644 index c4ad518e43..0000000000 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/TopicConfigCache.java +++ /dev/null @@ -1,94 +0,0 @@ -/* - * 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.connector; - -import com.google.common.cache.CacheBuilder; -import com.google.common.cache.LoadingCache; -import java.util.Optional; -import java.util.concurrent.ThreadPoolExecutor; -import java.util.concurrent.TimeUnit; -import org.apache.rocketmq.client.exception.MQClientException; -import org.apache.rocketmq.common.constant.LoggerName; -import org.apache.rocketmq.common.protocol.ResponseCode; -import org.apache.rocketmq.common.protocol.route.BrokerData; -import org.apache.rocketmq.common.statictopic.TopicConfigAndQueueMapping; -import org.apache.rocketmq.common.thread.ThreadPoolMonitor; -import org.apache.rocketmq.logging.InternalLogger; -import org.apache.rocketmq.logging.InternalLoggerFactory; -import org.apache.rocketmq.proxy.common.AbstractCacheLoader; -import org.apache.rocketmq.proxy.config.ConfigurationManager; -import org.apache.rocketmq.proxy.config.ProxyConfig; -import org.apache.rocketmq.proxy.connector.route.MessageQueueWrapper; -import org.apache.rocketmq.proxy.connector.route.TopicRouteCache; - -public class TopicConfigCache { - private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME); - - private final TopicRouteCache topicRouteCache; - private final ThreadPoolExecutor cacheRefreshExecutor; - private final LoadingCache topicConfigCache; - - private final DefaultForwardClient defaultClient; - - public TopicConfigCache(TopicRouteCache topicRouteCache, DefaultForwardClient client) { - this.topicRouteCache = topicRouteCache; - this.defaultClient = client; - - ProxyConfig config = ConfigurationManager.getProxyConfig(); - this.cacheRefreshExecutor = ThreadPoolMonitor.createAndMonitor( - config.getTopicConfigThreadPoolNums(), - config.getTopicConfigThreadPoolNums(), - 1000 * 60, - TimeUnit.MILLISECONDS, - "TopicConfigCacheRefresh", - config.getTopicConfigThreadPoolQueueCapacity() - ); - this.topicConfigCache = CacheBuilder.newBuilder() - .maximumSize(config.getTopicConfigCacheMaxNum()) - .refreshAfterWrite(config.getTopicConfigCacheExpiredInSeconds(), TimeUnit.SECONDS) - .build(new TopicConfigCacheLoader()); - } - - public TopicConfigAndQueueMapping getTopicConfigAndQueueMapping(String topic) throws Exception { - return topicConfigCache.get(topic); - } - - protected class TopicConfigCacheLoader extends AbstractCacheLoader { - - public TopicConfigCacheLoader() { - super(cacheRefreshExecutor); - } - - @Override - protected TopicConfigAndQueueMapping getDirectly(String topic) throws Exception { - MessageQueueWrapper messageQueueWrapper = topicRouteCache.getMessageQueue(topic); - Optional brokerDataOptional = messageQueueWrapper.getTopicRouteData().getBrokerDatas().stream().findAny(); - if (!brokerDataOptional.isPresent()) { - throw new MQClientException(ResponseCode.TOPIC_NOT_EXIST, - "No topic route info in name server for the topic: " + topic); - } - - String brokerAddress = brokerDataOptional.get().selectBrokerAddr(); - return defaultClient.getTopicConfig(brokerAddress, topic); - } - - @Override - protected void onErr(String key, Exception e) { - log.error("load topic config failed. topic:{}", key, e); - } - } -} 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 de6eda54a2..74b0052b03 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 @@ -32,12 +32,14 @@ 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; import org.apache.rocketmq.proxy.service.message.ClusterMessageService; import org.apache.rocketmq.proxy.service.message.MessageService; +import org.apache.rocketmq.proxy.service.metadata.ClusterMetadataService; +import org.apache.rocketmq.proxy.service.metadata.MetadataService; import org.apache.rocketmq.proxy.service.mqclient.DoNothingClientRemotingProcessor; import org.apache.rocketmq.proxy.service.relay.ClusterProxyRelayService; import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; +import org.apache.rocketmq.proxy.service.mqclient.MQClientAPIFactory; import org.apache.rocketmq.proxy.service.route.ClusterTopicRouteService; import org.apache.rocketmq.proxy.service.route.TopicRouteService; import org.apache.rocketmq.proxy.service.transaction.ClusterTransactionService; @@ -53,9 +55,11 @@ 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 ScheduledExecutorService scheduledExecutorService; - private final MQClientAPIFactory mqClientAPIFactory; + private final MQClientAPIFactory messagingClientAPIFactory; + private final MQClientAPIFactory operationClientAPIFactory; public ClusterServiceManager(RPCHook rpcHook) { this.scheduledExecutorService = Executors.newScheduledThreadPool(3); @@ -63,17 +67,25 @@ public class ClusterServiceManager extends AbstractStartAndShutdown implements S this.consumerManager = new ConsumerManager(new ConsumerIdsChangeListenerImpl()); ProxyConfig proxyConfig = ConfigurationManager.getProxyConfig(); - this.mqClientAPIFactory = new MQClientAPIFactory( + this.messagingClientAPIFactory = new MQClientAPIFactory( "CLUSTER_MQ_CLIENT_", proxyConfig.getRocketmqMQClientNum(), new DoNothingClientRemotingProcessor(null), rpcHook, scheduledExecutorService); + this.operationClientAPIFactory = new MQClientAPIFactory( + "TopicRouteServiceClient_", + 1, + new DoNothingClientRemotingProcessor(null), + rpcHook, + this.scheduledExecutorService + ); - this.topicRouteService = new ClusterTopicRouteService(rpcHook); - this.messageService = new ClusterMessageService(this.topicRouteService, this.mqClientAPIFactory); + this.topicRouteService = new ClusterTopicRouteService(operationClientAPIFactory); + this.messageService = new ClusterMessageService(this.topicRouteService, this.messagingClientAPIFactory); this.clusterTransactionService = new ClusterTransactionService(this.topicRouteService, this.producerManager, rpcHook); this.proxyRelayService = new ClusterProxyRelayService(); + this.metadataService = new ClusterMetadataService(topicRouteService, operationClientAPIFactory); this.init(); } @@ -91,7 +103,8 @@ public class ClusterServiceManager extends AbstractStartAndShutdown implements S }, 1000 * 10, 1000 * 10, TimeUnit.MILLISECONDS); this.appendShutdown(scheduledExecutorService::shutdown); - this.appendStartAndShutdown(this.mqClientAPIFactory); + this.appendStartAndShutdown(this.messagingClientAPIFactory); + this.appendStartAndShutdown(this.operationClientAPIFactory); this.appendStartAndShutdown(this.topicRouteService); this.appendStartAndShutdown(this.clusterTransactionService); } @@ -126,6 +139,11 @@ public class ClusterServiceManager extends AbstractStartAndShutdown implements S return this.proxyRelayService; } + @Override + public MetadataService getMetadataService() { + return this.metadataService; + } + protected static class ConsumerIdsChangeListenerImpl implements ConsumerIdsChangeListener { @Override 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 36a798d273..c55d0acb87 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 @@ -16,6 +16,7 @@ */ package org.apache.rocketmq.proxy.service; +import java.util.concurrent.Executors; import org.apache.rocketmq.broker.BrokerController; import org.apache.rocketmq.broker.client.ConsumerManager; import org.apache.rocketmq.broker.client.ProducerManager; @@ -24,6 +25,10 @@ import org.apache.rocketmq.proxy.service.message.LocalMessageService; import org.apache.rocketmq.proxy.service.message.MessageService; import org.apache.rocketmq.proxy.service.relay.LocalProxyRelayService; import org.apache.rocketmq.proxy.service.relay.ProxyRelayService; +import org.apache.rocketmq.proxy.service.metadata.LocalMetadataService; +import org.apache.rocketmq.proxy.service.metadata.MetadataService; +import org.apache.rocketmq.proxy.service.mqclient.DoNothingClientRemotingProcessor; +import org.apache.rocketmq.proxy.service.mqclient.MQClientAPIFactory; import org.apache.rocketmq.proxy.service.route.LocalTopicRouteService; import org.apache.rocketmq.proxy.service.route.TopicRouteService; import org.apache.rocketmq.proxy.service.transaction.LocalTransactionService; @@ -37,18 +42,29 @@ public class LocalServiceManager extends AbstractStartAndShutdown implements Ser private final MessageService messageService; private final TransactionService transactionService; private final ProxyRelayService proxyRelayService; + private final MetadataService metadataService; + + private final MQClientAPIFactory mqClientAPIFactory; public LocalServiceManager(BrokerController brokerController, RPCHook rpcHook) { this.brokerController = brokerController; this.messageService = new LocalMessageService(brokerController, rpcHook); - this.topicRouteService = new LocalTopicRouteService(brokerController, rpcHook); + this.mqClientAPIFactory = new MQClientAPIFactory( + "TopicRouteServiceClient_", + 1, + new DoNothingClientRemotingProcessor(null), + rpcHook, + Executors.newSingleThreadScheduledExecutor() + ); + this.topicRouteService = new LocalTopicRouteService(brokerController, mqClientAPIFactory); this.transactionService = new LocalTransactionService(); this.proxyRelayService = new LocalProxyRelayService(brokerController); - + this.metadataService = new LocalMetadataService(brokerController); this.init(); } protected void init() { + this.appendStartAndShutdown(this.mqClientAPIFactory); this.appendStartAndShutdown(this.topicRouteService); } @@ -82,4 +98,8 @@ public class LocalServiceManager extends AbstractStartAndShutdown implements Ser return this.proxyRelayService; } + @Override + public MetadataService getMetadataService() { + return this.metadataService; + } } 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 c80811f469..6a4f1cf371 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 @@ -21,6 +21,7 @@ import org.apache.rocketmq.broker.client.ProducerManager; 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.metadata.MetadataService; import org.apache.rocketmq.proxy.service.route.TopicRouteService; import org.apache.rocketmq.proxy.service.transaction.TransactionService; @@ -36,4 +37,6 @@ public interface ServiceManager extends StartAndShutdown { TransactionService getTransactionService(); ProxyRelayService getProxyRelayService(); + + MetadataService getMetadataService(); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/AbstractMetadataService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/AbstractMetadataService.java new file mode 100644 index 0000000000..e6ca6bdfc1 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/AbstractMetadataService.java @@ -0,0 +1,69 @@ +/* + * 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.metadata; + +import java.util.concurrent.ThreadPoolExecutor; +import java.util.concurrent.TimeUnit; +import org.apache.rocketmq.common.attribute.TopicMessageType; +import org.apache.rocketmq.common.constant.LoggerName; +import org.apache.rocketmq.common.statictopic.TopicConfigAndQueueMapping; +import org.apache.rocketmq.common.thread.ThreadPoolMonitor; +import org.apache.rocketmq.logging.InternalLogger; +import org.apache.rocketmq.logging.InternalLoggerFactory; +import org.apache.rocketmq.proxy.common.AbstractCacheLoader; +import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown; +import org.apache.rocketmq.proxy.config.ConfigurationManager; +import org.apache.rocketmq.proxy.config.ProxyConfig; + +public abstract class AbstractMetadataService extends AbstractStartAndShutdown implements MetadataService { + protected static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME); + protected final ThreadPoolExecutor cacheRefreshExecutor; + + public AbstractMetadataService() { + ProxyConfig config = ConfigurationManager.getProxyConfig(); + this.cacheRefreshExecutor = ThreadPoolMonitor.createAndMonitor( + config.getTopicConfigThreadPoolNums(), + config.getTopicConfigThreadPoolNums(), + 1000 * 60, + TimeUnit.MILLISECONDS, + "MetadataCacheRefresh", + config.getTopicConfigThreadPoolQueueCapacity() + ); + } + + public abstract TopicMessageType getTopicMessageType(String topic); + + protected abstract class AbstractTopicConfigCacheLoader extends AbstractCacheLoader { + + public AbstractTopicConfigCacheLoader() { + super(cacheRefreshExecutor); + } + + protected abstract TopicConfigAndQueueMapping loadTopicConfig(String topic) throws Exception; + + @Override + public TopicConfigAndQueueMapping getDirectly(String topic) throws Exception { + return loadTopicConfig(topic); + } + + @Override + protected void onErr(String key, Exception e) { + log.error("load topic config failed. topic:{}", key, e); + } + } +} 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 new file mode 100644 index 0000000000..52acefe127 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/ClusterMetadataService.java @@ -0,0 +1,85 @@ +/* + * 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.metadata; + +import com.google.common.cache.CacheBuilder; +import com.google.common.cache.LoadingCache; +import java.util.Optional; +import java.util.concurrent.TimeUnit; +import org.apache.rocketmq.client.exception.MQClientException; +import org.apache.rocketmq.common.attribute.TopicMessageType; +import org.apache.rocketmq.common.protocol.route.BrokerData; +import org.apache.rocketmq.common.statictopic.TopicConfigAndQueueMapping; +import org.apache.rocketmq.proxy.config.ConfigurationManager; +import org.apache.rocketmq.proxy.config.ProxyConfig; +import org.apache.rocketmq.proxy.service.mqclient.MQClientAPIFactory; +import org.apache.rocketmq.proxy.service.route.TopicRouteHelper; +import org.apache.rocketmq.proxy.service.route.TopicRouteService; + +public class ClusterMetadataService extends AbstractMetadataService { + private final LoadingCache topicCache; + private TopicRouteService topicRouteService; + private final static TopicConfigAndQueueMapping EMPTY_TOPIC_CONFIG = new TopicConfigAndQueueMapping(); + + public ClusterMetadataService(TopicRouteService topicRouteService, MQClientAPIFactory mqClientAPIFactory) { + this.topicRouteService = topicRouteService; + ProxyConfig config = ConfigurationManager.getProxyConfig(); + this.topicCache = CacheBuilder.newBuilder() + .maximumSize(config.getTopicRouteCacheMaxNum()) + .refreshAfterWrite(config.getTopicRouteCacheExpiredInSeconds(), TimeUnit.SECONDS) + .build(new ClusterTopicConfigCacheLoader(mqClientAPIFactory)); + } + + @Override public TopicMessageType getTopicMessageType(String topic) { + TopicConfigAndQueueMapping topicConfigAndQueueMapping; + try { + topicConfigAndQueueMapping = topicCache.get(topic); + } catch (Exception e) { + return TopicMessageType.UNSPECIFIED; + } + if (topicConfigAndQueueMapping.equals(EMPTY_TOPIC_CONFIG)) { + return TopicMessageType.UNSPECIFIED; + } + return topicConfigAndQueueMapping.getTopicMessageType(); + } + + protected class ClusterTopicConfigCacheLoader extends AbstractTopicConfigCacheLoader { + private final MQClientAPIFactory mqClientAPIFactory; + + public ClusterTopicConfigCacheLoader(MQClientAPIFactory mqClientAPIFactory) { + this.mqClientAPIFactory = mqClientAPIFactory; + } + + @Override protected TopicConfigAndQueueMapping loadTopicConfig(String topic) throws Exception { + try { + Optional brokerDataOptional = topicRouteService.getAllMessageQueueView(topic).getTopicRouteData().getBrokerDatas().stream().findAny(); + if (brokerDataOptional.isPresent()) { + String brokerAddress = topicRouteService.getBrokerAddr(brokerDataOptional.get().getBrokerName()); + return this.mqClientAPIFactory.getClient().getTopicConfig(brokerAddress, topic, 1000L); + } + return EMPTY_TOPIC_CONFIG; + } catch (MQClientException e) { + if (TopicRouteHelper.isTopicNotExistError(e)) { + log.warn("topic is not exist", e); + return EMPTY_TOPIC_CONFIG; + } + throw e; + } + } + } +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/LocalMetadataService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/LocalMetadataService.java new file mode 100644 index 0000000000..7bf0ece856 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/LocalMetadataService.java @@ -0,0 +1,68 @@ +/* + * 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.metadata; + +import com.google.common.cache.CacheBuilder; +import com.google.common.cache.LoadingCache; +import java.util.concurrent.TimeUnit; +import org.apache.rocketmq.broker.BrokerController; +import org.apache.rocketmq.common.TopicConfig; +import org.apache.rocketmq.common.attribute.TopicMessageType; +import org.apache.rocketmq.common.statictopic.TopicConfigAndQueueMapping; +import org.apache.rocketmq.proxy.config.ConfigurationManager; +import org.apache.rocketmq.proxy.config.ProxyConfig; + +public class LocalMetadataService extends AbstractMetadataService { + private final LoadingCache topicCache; + + public LocalMetadataService(BrokerController brokerController) { + ProxyConfig config = ConfigurationManager.getProxyConfig(); + + this.topicCache = CacheBuilder.newBuilder() + .maximumSize(config.getTopicRouteCacheMaxNum()) + .refreshAfterWrite(config.getTopicRouteCacheExpiredInSeconds(), TimeUnit.SECONDS) + .build(new LocalTopicConfigCacheLoader(brokerController)); + } + + @Override + public TopicMessageType getTopicMessageType(String topic) { + try { + TopicConfigAndQueueMapping topicConfigAndQueueMapping = topicCache.get(topic); + if (topicConfigAndQueueMapping == null) { + return TopicMessageType.UNSPECIFIED; + } + return topicConfigAndQueueMapping.getTopicMessageType(); + } catch (Exception e) { + log.error("getTopicMessageType error", topic, e); + return TopicMessageType.UNSPECIFIED; + } + } + + protected class LocalTopicConfigCacheLoader extends AbstractTopicConfigCacheLoader { + private final BrokerController brokerController; + + public LocalTopicConfigCacheLoader(BrokerController brokerController) { + this.brokerController = brokerController; + } + + @Override protected TopicConfigAndQueueMapping loadTopicConfig(String topic) throws Exception { + TopicConfig topicConfig = brokerController.getTopicConfigManager().selectTopicConfig(topic); + return new TopicConfigAndQueueMapping(topicConfig, null); + } + } +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/MetadataService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/MetadataService.java new file mode 100644 index 0000000000..a87ec1484f --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/metadata/MetadataService.java @@ -0,0 +1,24 @@ +/* + * 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.metadata; + +import org.apache.rocketmq.common.attribute.TopicMessageType; + +public interface MetadataService { + TopicMessageType getTopicMessageType(String topic); +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/ClusterTopicRouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/ClusterTopicRouteService.java index de058db13c..94d7dbdd5f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/ClusterTopicRouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/ClusterTopicRouteService.java @@ -21,12 +21,12 @@ import org.apache.rocketmq.common.MixAll; import org.apache.rocketmq.common.protocol.route.BrokerData; import org.apache.rocketmq.common.protocol.route.TopicRouteData; import org.apache.rocketmq.proxy.common.Address; -import org.apache.rocketmq.remoting.RPCHook; +import org.apache.rocketmq.proxy.service.mqclient.MQClientAPIFactory; public class ClusterTopicRouteService extends TopicRouteService { - public ClusterTopicRouteService(RPCHook rpcHook) { - super(rpcHook); + public ClusterTopicRouteService(MQClientAPIFactory mqClientAPIFactory) { + super(mqClientAPIFactory); } @Override diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/LocalTopicRouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/LocalTopicRouteService.java index 87d9078997..5da75cea08 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/LocalTopicRouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/LocalTopicRouteService.java @@ -29,7 +29,7 @@ import org.apache.rocketmq.common.protocol.route.QueueData; import org.apache.rocketmq.common.protocol.route.TopicRouteData; import org.apache.rocketmq.proxy.common.Address; import org.apache.rocketmq.proxy.config.ConfigurationManager; -import org.apache.rocketmq.remoting.RPCHook; +import org.apache.rocketmq.proxy.service.mqclient.MQClientAPIFactory; public class LocalTopicRouteService extends TopicRouteService { @@ -37,8 +37,8 @@ public class LocalTopicRouteService extends TopicRouteService { private final List brokerDataList; private final int grpcPort; - public LocalTopicRouteService(BrokerController brokerController, RPCHook rpcHook) { - super(rpcHook); + public LocalTopicRouteService(BrokerController brokerController, MQClientAPIFactory mqClientAPIFactory) { + super(mqClientAPIFactory); this.brokerController = brokerController; BrokerConfig brokerConfig = this.brokerController.getBrokerConfig(); HashMap brokerAddrs = new HashMap<>(); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/TopicRouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/TopicRouteService.java index aa7e5ec359..c800abfa1a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/TopicRouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/service/route/TopicRouteService.java @@ -37,9 +37,7 @@ import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown; import org.apache.rocketmq.proxy.common.Address; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.config.ProxyConfig; -import org.apache.rocketmq.proxy.service.mqclient.DoNothingClientRemotingProcessor; import org.apache.rocketmq.proxy.service.mqclient.MQClientAPIFactory; -import org.apache.rocketmq.remoting.RPCHook; public abstract class TopicRouteService extends AbstractStartAndShutdown { private static final InternalLogger log = InternalLoggerFactory.getLogger(LoggerName.PROXY_LOGGER_NAME); @@ -50,7 +48,7 @@ public abstract class TopicRouteService extends AbstractStartAndShutdown { private final ScheduledExecutorService scheduledExecutorService; private final ThreadPoolExecutor cacheRefreshExecutor; - public TopicRouteService(RPCHook rpcHook) { + public TopicRouteService(MQClientAPIFactory mqClientAPIFactory) { ProxyConfig config = ConfigurationManager.getProxyConfig(); this.scheduledExecutorService = Executors.newSingleThreadScheduledExecutor( @@ -64,13 +62,7 @@ public abstract class TopicRouteService extends AbstractStartAndShutdown { "TopicRouteCacheRefresh", config.getTopicRouteServiceThreadPoolQueueCapacity() ); - this.mqClientAPIFactory = new MQClientAPIFactory( - "TopicRouteServiceClient_", - 1, - new DoNothingClientRemotingProcessor(null), - rpcHook, - this.scheduledExecutorService - ); + this.mqClientAPIFactory = mqClientAPIFactory; this.topicCache = CacheBuilder.newBuilder() .maximumSize(config.getTopicRouteServiceCacheMaxNum()) .refreshAfterWrite(config.getTopicRouteServiceCacheExpiredInSeconds(), TimeUnit.SECONDS)