From 3b437435b479d614bfcbaebc38d10c2de5abbc2b Mon Sep 17 00:00:00 2001 From: "Jixiang.jjx" Date: Tue, 15 Mar 2022 16:46:29 +0800 Subject: [PATCH] [ISSUE #3949] do some renaming work. --- .../apache/rocketmq/proxy/ProxyStartup.java | 4 +- .../proxy/channel/ChannelManager.java | 5 +- ...Client.java => AbstractForwardClient.java} | 13 +- .../rocketmq/proxy/client/ClientManager.java | 67 ----- ...tClient.java => DefaultForwardClient.java} | 19 +- .../proxy/client/ForwardClientManager.java | 71 ++++++ ...oducerClient.java => ForwardProducer.java} | 13 +- ...erClient.java => ForwardReadConsumer.java} | 11 +- ...rClient.java => ForwardWriteConsumer.java} | 11 +- .../proxy/client/TopicRouteCache.java | 26 +- .../AbstractMQClientFactory.java} | 10 +- .../ForwardClientFactory.java} | 28 +-- .../MQClientFactory.java} | 6 +- .../MQClientFactoryImpl.java} | 6 +- .../TransactionalProducerFactory.java} | 11 +- .../DoNothingClientRemotingProcessor.java | 3 +- .../ProxyClientRemotingProcessor.java | 18 +- .../client/route/AddressableMessageQueue.java | 80 ------ .../client/route/MessageQueueSelector.java | 229 ++++++++++++++++++ .../client/route/MessageQueueWrapper.java | 24 +- .../client/route/SelectableMessageQueue.java | 219 +++-------------- .../client/transaction/TransactionId.java | 18 +- .../TransactionStateCheckRequest.java | 10 +- .../transaction/TransactionStateChecker.java | 1 - .../proxy/configuration/ProxyConfig.java | 90 +++---- .../grpc/service/ClusterGrpcService.java | 6 +- .../grpc/service/cluster/BaseService.java | 6 +- .../grpc/service/cluster/ConsumerService.java | 4 +- .../grpc/service/cluster/ProducerService.java | 26 +- .../grpc/service/cluster/RouteService.java | 16 +- .../proxy/client/ClientManagerTest.java | 20 +- .../grpc/service/cluster/BaseServiceTest.java | 28 +-- .../service/cluster/ProducerServiceTest.java | 26 +- 33 files changed, 577 insertions(+), 548 deletions(-) rename proxy/src/main/java/org/apache/rocketmq/proxy/client/{BaseClient.java => AbstractForwardClient.java} (76%) delete mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/client/ClientManager.java rename proxy/src/main/java/org/apache/rocketmq/proxy/client/{DefaultClient.java => DefaultForwardClient.java} (77%) create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/client/ForwardClientManager.java rename proxy/src/main/java/org/apache/rocketmq/proxy/client/{ProducerClient.java => ForwardProducer.java} (84%) rename proxy/src/main/java/org/apache/rocketmq/proxy/client/{ReadConsumerClient.java => ForwardReadConsumer.java} (84%) rename proxy/src/main/java/org/apache/rocketmq/proxy/client/{WriteConsumerClient.java => ForwardWriteConsumer.java} (86%) rename proxy/src/main/java/org/apache/rocketmq/proxy/client/{mqconstructor/AbstractRocketMQClientConstructor.java => factory/AbstractMQClientFactory.java} (86%) rename proxy/src/main/java/org/apache/rocketmq/proxy/client/{ClientFactory.java => factory/ForwardClientFactory.java} (69%) rename proxy/src/main/java/org/apache/rocketmq/proxy/client/{mqconstructor/RocketMQClientConstructor.java => factory/MQClientFactory.java} (89%) rename proxy/src/main/java/org/apache/rocketmq/proxy/client/{mqconstructor/MQClientAPIConstructor.java => factory/MQClientFactoryImpl.java} (89%) rename proxy/src/main/java/org/apache/rocketmq/proxy/client/{mqconstructor/TransactionClientConstructor.java => factory/TransactionalProducerFactory.java} (74%) delete mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/client/route/AddressableMessageQueue.java create mode 100644 proxy/src/main/java/org/apache/rocketmq/proxy/client/route/MessageQueueSelector.java diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java b/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java index e70a227cbd..8379fcf66c 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java @@ -108,7 +108,7 @@ public class ProxyStartup { private static void initThreadPoolMonitor() { ThreadPoolMonitor.init(); ProxyConfig config = ConfigurationManager.getProxyConfig(); - ThreadPoolMonitor.config(config.isEnablePrintJstack(), config.getPrintJstackPeriodMillis()); + ThreadPoolMonitor.config(config.isEnablePrintJstack(), config.getPrintJstackInMillis()); } private static void initLogger() throws JoranException { @@ -120,6 +120,6 @@ public class ProxyStartup { lc.reset(); //https://logback.qos.ch/manual/configuration.html lc.setPackagingDataEnabled(false); - configurator.doConfigure(ConfigurationManager.getProxyHome() + "/conf/logback.xml"); + configurator.doConfigure(ConfigurationManager.getProxyHome() + "/conf/logback_proxy.xml"); } } \ No newline at end of file diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/channel/ChannelManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/channel/ChannelManager.java index 2c9ff07d9a..d8666fa5b5 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/channel/ChannelManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/channel/ChannelManager.java @@ -63,7 +63,7 @@ public class ChannelManager { .get(InterceptorConstants.REMOTE_ADDRESS); final String localAddress = InterceptorConstants.METADATA.get(Context.current()) .get(InterceptorConstants.LOCAL_ADDRESS); - return new SimpleChannel(null, clientHost, localAddress, ConfigurationManager.getProxyConfig().getExpiredChannelTimeSec()); + return new SimpleChannel(null, clientHost, localAddress, ConfigurationManager.getProxyConfig().getChannelExpiredInSeconds()); } /** @@ -71,8 +71,7 @@ public class ChannelManager { */ public void scanAndCleanChannels() { try { - Iterator> iterator = clientIdChannelMap.entrySet() - .iterator(); + Iterator> iterator = clientIdChannelMap.entrySet().iterator(); while (iterator.hasNext()) { Map.Entry entry = iterator.next(); if (!entry.getValue() diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/BaseClient.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/AbstractForwardClient.java similarity index 76% rename from proxy/src/main/java/org/apache/rocketmq/proxy/client/BaseClient.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/client/AbstractForwardClient.java index 270e1ae825..480554d6b8 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/BaseClient.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/AbstractForwardClient.java @@ -18,20 +18,21 @@ package org.apache.rocketmq.proxy.client; import java.util.concurrent.ThreadLocalRandom; import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; +import org.apache.rocketmq.proxy.client.factory.ForwardClientFactory; import org.apache.rocketmq.proxy.common.StartAndShutdown; -public abstract class BaseClient implements StartAndShutdown { +public abstract class AbstractForwardClient implements StartAndShutdown { - private final ClientFactory clientFactory; + private final ForwardClientFactory forwardClientFactory; private MQClientAPIExtImpl[] clients; - public BaseClient(ClientFactory clientFactory) { - this.clientFactory = clientFactory; + public AbstractForwardClient(ForwardClientFactory forwardClientFactory) { + this.forwardClientFactory = forwardClientFactory; } protected abstract int getClientNum(); - protected abstract MQClientAPIExtImpl createNewClient(ClientFactory clientFactory, String name); + protected abstract MQClientAPIExtImpl createNewClient(ForwardClientFactory forwardClientFactory, String name); protected abstract String getNamePrefix(); @@ -48,7 +49,7 @@ public abstract class BaseClient implements StartAndShutdown { this.clients = new MQClientAPIExtImpl[clientCount]; for (int i = 0; i < clientCount; i++) { String name = getNamePrefix() + "N_" + i; - clients[i] = createNewClient(clientFactory, name); + clients[i] = createNewClient(forwardClientFactory, name); } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/ClientManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/ClientManager.java deleted file mode 100644 index 5c571620cd..0000000000 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/ClientManager.java +++ /dev/null @@ -1,67 +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.client; - -import org.apache.rocketmq.proxy.client.transaction.TransactionStateChecker; -import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown; - -public class ClientManager extends AbstractStartAndShutdown { - - private final ClientFactory clientFactory; - private final DefaultClient defaultClient; - private final ProducerClient producerClient; - private final ReadConsumerClient readConsumerClient; - private final WriteConsumerClient writeConsumerClient; - - private final TopicRouteCache topicRouteCache; - - public ClientManager(TransactionStateChecker transactionStateChecker) { - this.clientFactory = new ClientFactory(transactionStateChecker); - this.defaultClient = new DefaultClient(this.clientFactory); - this.producerClient = new ProducerClient(this.clientFactory); - this.readConsumerClient = new ReadConsumerClient(this.clientFactory); - this.writeConsumerClient = new WriteConsumerClient(this.clientFactory); - - this.topicRouteCache = new TopicRouteCache(this.defaultClient); - - this.appendStartAndShutdown(this.clientFactory); - this.appendStartAndShutdown(this.defaultClient); - this.appendStartAndShutdown(this.producerClient); - this.appendStartAndShutdown(this.readConsumerClient); - this.appendStartAndShutdown(this.writeConsumerClient); - } - - public DefaultClient getDefaultClient() { - return defaultClient; - } - - public ProducerClient getProducerClient() { - return producerClient; - } - - public ReadConsumerClient getReadConsumerClient() { - return readConsumerClient; - } - - public WriteConsumerClient getWriteConsumerClient() { - return writeConsumerClient; - } - - public TopicRouteCache getTopicRouteCache() { - return topicRouteCache; - } -} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/DefaultClient.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/DefaultForwardClient.java similarity index 77% rename from proxy/src/main/java/org/apache/rocketmq/proxy/client/DefaultClient.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/client/DefaultForwardClient.java index 4ab7518eb4..de3d950160 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/DefaultClient.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/DefaultForwardClient.java @@ -22,25 +22,25 @@ import org.apache.rocketmq.client.exception.MQClientException; import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; import org.apache.rocketmq.common.protocol.header.GetConsumerListByGroupRequestHeader; import org.apache.rocketmq.common.protocol.route.TopicRouteData; +import org.apache.rocketmq.proxy.client.factory.ForwardClientFactory; import org.apache.rocketmq.proxy.configuration.ConfigurationManager; import org.apache.rocketmq.remoting.exception.RemotingException; -public class DefaultClient extends BaseClient { - +public class DefaultForwardClient extends AbstractForwardClient { private static final String CID_PREFIX = "CID_RMQ_PROXY_DEFAULT_"; - public DefaultClient(ClientFactory clientFactory) { + public DefaultForwardClient(ForwardClientFactory clientFactory) { super(clientFactory); } @Override protected int getClientNum() { - return ConfigurationManager.getProxyConfig().getDefaultClientNum(); + return ConfigurationManager.getProxyConfig().getDefaultForwardClientNum(); } @Override - protected MQClientAPIExtImpl createNewClient(ClientFactory clientFactory, String name) { - double workerFactor = ConfigurationManager.getProxyConfig().getDefaultClientWorkerFactor(); + protected MQClientAPIExtImpl createNewClient(ForwardClientFactory clientFactory, String name) { + double workerFactor = ConfigurationManager.getProxyConfig().getDefaultForwardClientWorkerFactor(); final int threadCount = (int) Math.ceil(Runtime.getRuntime().availableProcessors() * workerFactor); return clientFactory.getMQClient(name, threadCount); @@ -51,12 +51,13 @@ public class DefaultClient extends BaseClient { return CID_PREFIX; } - public CompletableFuture> getConsumerListByGroup(String brokerAddr, GetConsumerListByGroupRequestHeader requestHeader, - long timeoutMillis) { + public CompletableFuture> getConsumerListByGroup( + String brokerAddr, GetConsumerListByGroupRequestHeader requestHeader, long timeoutMillis) { return getClient().getConsumerListByGroup(brokerAddr, requestHeader, timeoutMillis); } - public TopicRouteData getTopicRouteInfoFromNameServer(String topic, long timeoutMillis) throws RemotingException, InterruptedException, MQClientException { + public TopicRouteData getTopicRouteInfoFromNameServer(String topic, long timeoutMillis) + throws RemotingException, InterruptedException, MQClientException { return getClient().getTopicRouteInfoFromNameServer(topic, timeoutMillis); } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/ForwardClientManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/ForwardClientManager.java new file mode 100644 index 0000000000..85dc298a8b --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/ForwardClientManager.java @@ -0,0 +1,71 @@ +/* + * 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.client; + +import org.apache.rocketmq.proxy.client.factory.ForwardClientFactory; +import org.apache.rocketmq.proxy.client.transaction.TransactionStateChecker; +import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown; + +public class ForwardClientManager extends AbstractStartAndShutdown { + private final ForwardClientFactory forwardClientFactory; + private final DefaultForwardClient defaultForwardClient; + private final ForwardProducer forwardProducer; + private final ForwardReadConsumer forwardReadConsumer; + private final ForwardWriteConsumer forwardWriteConsumer; + + private final TopicRouteCache topicRouteCache; + + public ForwardClientManager(TransactionStateChecker transactionStateChecker) { + this.forwardClientFactory = new ForwardClientFactory(transactionStateChecker); + this.defaultForwardClient = new DefaultForwardClient(this.forwardClientFactory); + this.forwardProducer = new ForwardProducer(this.forwardClientFactory); + this.forwardReadConsumer = new ForwardReadConsumer(this.forwardClientFactory); + this.forwardWriteConsumer = new ForwardWriteConsumer(this.forwardClientFactory); + + this.topicRouteCache = new TopicRouteCache(this.defaultForwardClient); + + this.appendStartAndShutdown(this.forwardClientFactory); + this.appendStartAndShutdown(this.defaultForwardClient); + this.appendStartAndShutdown(this.forwardProducer); + this.appendStartAndShutdown(this.forwardReadConsumer); + this.appendStartAndShutdown(this.forwardWriteConsumer); + } + + public ForwardClientFactory getForwardClientFactory() { + return forwardClientFactory; + } + + public DefaultForwardClient getDefaultForwardClient() { + return defaultForwardClient; + } + + public ForwardProducer getForwardProducer() { + return forwardProducer; + } + + public ForwardReadConsumer getForwardReadConsumer() { + return forwardReadConsumer; + } + + public ForwardWriteConsumer getForwardWriteConsumer() { + return forwardWriteConsumer; + } + + public TopicRouteCache getTopicRouteCache() { + return topicRouteCache; + } +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/ProducerClient.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/ForwardProducer.java similarity index 84% rename from proxy/src/main/java/org/apache/rocketmq/proxy/client/ProducerClient.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/client/ForwardProducer.java index 668cb7fce8..8254b83fb1 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/ProducerClient.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/ForwardProducer.java @@ -23,28 +23,29 @@ import org.apache.rocketmq.common.message.Message; import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader; import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader; import org.apache.rocketmq.common.protocol.heartbeat.HeartbeatData; +import org.apache.rocketmq.proxy.client.factory.ForwardClientFactory; import org.apache.rocketmq.proxy.configuration.ConfigurationManager; import org.apache.rocketmq.remoting.protocol.RemotingCommand; -public class ProducerClient extends BaseClient { +public class ForwardProducer extends AbstractForwardClient { private static final String PID_PREFIX = "PID_RMQ_PROXY_PUBLISH_MESSAGE_"; - public ProducerClient(ClientFactory clientFactory) { + public ForwardProducer(ForwardClientFactory clientFactory) { super(clientFactory); } @Override protected int getClientNum() { - return ConfigurationManager.getProxyConfig().getProducerClientNum(); + return ConfigurationManager.getProxyConfig().getForwardProducerNum(); } @Override - protected MQClientAPIExtImpl createNewClient(ClientFactory clientFactory, String name) { - double sendClientWorkerFactor = ConfigurationManager.getProxyConfig().getProducerClientWorkerFactor(); + protected MQClientAPIExtImpl createNewClient(ForwardClientFactory clientFactory, String name) { + double sendClientWorkerFactor = ConfigurationManager.getProxyConfig().getForwardProducerWorkerFactor(); final int threadCount = (int) Math.ceil(Runtime.getRuntime().availableProcessors() * sendClientWorkerFactor); - return clientFactory.getTransactionClient(name, threadCount); + return clientFactory.getTransactionalProducer(name, threadCount); } @Override diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/ReadConsumerClient.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/ForwardReadConsumer.java similarity index 84% rename from proxy/src/main/java/org/apache/rocketmq/proxy/client/ReadConsumerClient.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/client/ForwardReadConsumer.java index 60516f9534..f3672ea370 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/ReadConsumerClient.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/ForwardReadConsumer.java @@ -22,24 +22,25 @@ import org.apache.rocketmq.client.consumer.PullResult; import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader; import org.apache.rocketmq.common.protocol.header.PullMessageRequestHeader; +import org.apache.rocketmq.proxy.client.factory.ForwardClientFactory; import org.apache.rocketmq.proxy.configuration.ConfigurationManager; -public class ReadConsumerClient extends BaseClient { +public class ForwardReadConsumer extends AbstractForwardClient { private static final String CID_PREFIX = "CID_RMQ_PROXY_CONSUME_MESSAGE_"; - public ReadConsumerClient(ClientFactory clientFactory) { + public ForwardReadConsumer(ForwardClientFactory clientFactory) { super(clientFactory); } @Override protected int getClientNum() { - return ConfigurationManager.getProxyConfig().getConsumerClientNum(); + return ConfigurationManager.getProxyConfig().getForwardConsumerNum(); } @Override - protected MQClientAPIExtImpl createNewClient(ClientFactory clientFactory, String name) { - double workerFactor = ConfigurationManager.getProxyConfig().getConsumerClientWorkerFactor(); + protected MQClientAPIExtImpl createNewClient(ForwardClientFactory clientFactory, String name) { + double workerFactor = ConfigurationManager.getProxyConfig().getForwardConsumerWorkerFactor(); final int threadCount = (int) Math.ceil(Runtime.getRuntime().availableProcessors() * workerFactor); return clientFactory.getMQClient(name, threadCount); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/WriteConsumerClient.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/ForwardWriteConsumer.java similarity index 86% rename from proxy/src/main/java/org/apache/rocketmq/proxy/client/WriteConsumerClient.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/client/ForwardWriteConsumer.java index ee04262163..26dedbbd0d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/WriteConsumerClient.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/ForwardWriteConsumer.java @@ -22,25 +22,26 @@ import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; import org.apache.rocketmq.common.protocol.header.AckMessageRequestHeader; import org.apache.rocketmq.common.protocol.header.ChangeInvisibleTimeRequestHeader; import org.apache.rocketmq.common.protocol.header.UpdateConsumerOffsetRequestHeader; +import org.apache.rocketmq.proxy.client.factory.ForwardClientFactory; import org.apache.rocketmq.proxy.configuration.ConfigurationManager; import org.apache.rocketmq.remoting.exception.RemotingException; -public class WriteConsumerClient extends BaseClient { +public class ForwardWriteConsumer extends AbstractForwardClient { private static final String CID_PREFIX = "CID_RMQ_PROXY_DELETE_MESSAGE_"; - public WriteConsumerClient(ClientFactory clientFactory) { + public ForwardWriteConsumer(ForwardClientFactory clientFactory) { super(clientFactory); } @Override protected int getClientNum() { - return ConfigurationManager.getProxyConfig().getConsumerClientNum(); + return ConfigurationManager.getProxyConfig().getForwardConsumerNum(); } @Override - protected MQClientAPIExtImpl createNewClient(ClientFactory clientFactory, String name) { - double workerFactor = ConfigurationManager.getProxyConfig().getConsumerClientWorkerFactor(); + protected MQClientAPIExtImpl createNewClient(ForwardClientFactory clientFactory, String name) { + double workerFactor = ConfigurationManager.getProxyConfig().getForwardConsumerWorkerFactor(); final int threadCount = (int) Math.ceil(Runtime.getRuntime().availableProcessors() * workerFactor); return clientFactory.getMQClient(name, threadCount); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/TopicRouteCache.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/TopicRouteCache.java index dec5496668..33091ce136 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/TopicRouteCache.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/TopicRouteCache.java @@ -26,7 +26,7 @@ import org.apache.rocketmq.client.exception.MQClientException; import org.apache.rocketmq.common.protocol.ResponseCode; import org.apache.rocketmq.common.protocol.route.TopicRouteData; import org.apache.rocketmq.common.thread.ThreadPoolMonitor; -import org.apache.rocketmq.proxy.client.route.AddressableMessageQueue; +import org.apache.rocketmq.proxy.client.route.SelectableMessageQueue; import org.apache.rocketmq.proxy.client.route.MessageQueueWrapper; import org.apache.rocketmq.proxy.common.RetainCacheLoader; import org.apache.rocketmq.proxy.common.RocketMQHelper; @@ -43,9 +43,9 @@ public class TopicRouteCache { private final LoadingCache topicCache; private final ThreadPoolExecutor cacheRefreshExecutor; - private final DefaultClient defaultClient; + private final DefaultForwardClient defaultClient; - public TopicRouteCache(DefaultClient defaultClient) { + public TopicRouteCache(DefaultForwardClient defaultClient) { ProxyConfig config = ConfigurationManager.getProxyConfig(); this.defaultClient = defaultClient; @@ -59,7 +59,7 @@ public class TopicRouteCache { ); this.topicCache = CacheBuilder.newBuilder() .maximumSize(config.getTopicRouteCacheMaxNum()) - .refreshAfterWrite(config.getTopicRouteCacheExpireSecond(), TimeUnit.SECONDS) + .refreshAfterWrite(config.getTopicRouteCacheExpiredInSeconds(), TimeUnit.SECONDS) .build(new TopicRouteCacheLoader()); } @@ -67,19 +67,19 @@ public class TopicRouteCache { return getCacheMessageQueueWrapper(this.topicCache, topicName); } - public AddressableMessageQueue selectOneWriteQueue(String topic, AddressableMessageQueue last) throws Exception { + public SelectableMessageQueue selectOneWriteQueue(String topic, SelectableMessageQueue last) throws Exception { if (last == null) { - return getMessageQueue(topic).getWrite().selectOne(false); + return getMessageQueue(topic).getWriteSelector().selectOne(false); } - return getMessageQueue(topic).getWrite().selectNextQueue(last); + return getMessageQueue(topic).getWriteSelector().selectNextQueue(last); } - public AddressableMessageQueue selectOneWriteQueue(String topic, String brokerName, int queueId) throws Exception { - return getMessageQueue(topic).getWrite().selectOne(brokerName, queueId); + public SelectableMessageQueue selectOneWriteQueue(String topic, String brokerName, int queueId) throws Exception { + return getMessageQueue(topic).getWriteSelector().selectOne(brokerName, queueId); } - public AddressableMessageQueue selectOneWriteQueueByKey(String topic, String shardingKey, AddressableMessageQueue last) throws Exception { - List writeQueues = getMessageQueue(topic).getWrite().getQueues(); + public SelectableMessageQueue selectOneWriteQueueByKey(String topic, String shardingKey, SelectableMessageQueue last) throws Exception { + List writeQueues = getMessageQueue(topic).getWriteSelector().getQueues(); int bucket = Hashing.consistentHash(shardingKey.hashCode(), writeQueues.size()); return writeQueues.get(bucket); } @@ -122,10 +122,10 @@ public class TopicRouteCache { log.info("load {} from namesrv. topic: {}, queue: {}", loaderName(), topic, tmp); return tmp; } - return MessageQueueWrapper.EMPTY_CACHED_QUEUE; + return MessageQueueWrapper.WRAPPED_EMPTY_QUEUE; } catch (Exception e) { if (RocketMQHelper.isTopicNotExistError(e)) { - return MessageQueueWrapper.EMPTY_CACHED_QUEUE; + return MessageQueueWrapper.WRAPPED_EMPTY_QUEUE; } throw e; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/mqconstructor/AbstractRocketMQClientConstructor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/factory/AbstractMQClientFactory.java similarity index 86% rename from proxy/src/main/java/org/apache/rocketmq/proxy/client/mqconstructor/AbstractRocketMQClientConstructor.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/client/factory/AbstractMQClientFactory.java index 4e530cacf8..35221594f1 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/mqconstructor/AbstractRocketMQClientConstructor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/factory/AbstractMQClientFactory.java @@ -14,7 +14,7 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.rocketmq.proxy.client.mqconstructor; +package org.apache.rocketmq.proxy.client.factory; import java.util.Map; import java.util.concurrent.ConcurrentHashMap; @@ -23,14 +23,14 @@ import org.apache.rocketmq.remoting.netty.NettyClientConfig; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -public abstract class AbstractRocketMQClientConstructor implements RocketMQClientConstructor { +public abstract class AbstractMQClientFactory implements MQClientFactory { - private static final Logger log = LoggerFactory.getLogger(AbstractRocketMQClientConstructor.class); + private static final Logger LOGGER = LoggerFactory.getLogger(AbstractMQClientFactory.class); protected Map cacheTable = new ConcurrentHashMap<>(); protected RPCHook rpcHook; - public AbstractRocketMQClientConstructor(RPCHook rpcHook) { + public AbstractMQClientFactory(RPCHook rpcHook) { this.rpcHook = rpcHook; } @@ -78,7 +78,7 @@ public abstract class AbstractRocketMQClientConstructor implements RocketMQCl try { this.shutdown(v); } catch (Exception e) { - log.warn("RocketMQClientConstructor shutdown all err.", e); + LOGGER.warn("RocketMQClientConstructor shutdown all err.", e); } }); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/ClientFactory.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/factory/ForwardClientFactory.java similarity index 69% rename from proxy/src/main/java/org/apache/rocketmq/proxy/client/ClientFactory.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/client/factory/ForwardClientFactory.java index 5e777eef30..55e709cb4d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/ClientFactory.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/factory/ForwardClientFactory.java @@ -14,37 +14,35 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.rocketmq.proxy.client; +package org.apache.rocketmq.proxy.client.factory; import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.client.ClientConfig; import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; import org.apache.rocketmq.common.MixAll; -import org.apache.rocketmq.proxy.client.mqconstructor.MQClientAPIConstructor; -import org.apache.rocketmq.proxy.client.mqconstructor.TransactionClientConstructor; import org.apache.rocketmq.proxy.client.transaction.TransactionStateChecker; import org.apache.rocketmq.proxy.common.StartAndShutdown; import org.apache.rocketmq.proxy.configuration.ConfigurationManager; import org.apache.rocketmq.remoting.RPCHook; -public class ClientFactory implements StartAndShutdown { +public class ForwardClientFactory implements StartAndShutdown { private RPCHook rpcHook = null; - private final MQClientAPIConstructor mqClientAPIConstructor; - private final TransactionClientConstructor transactionClientConstructor; + private final MQClientFactoryImpl mqClientFactory; + private final TransactionalProducerFactory transactionalProducerFactory; - public ClientFactory(TransactionStateChecker transactionStateChecker) { + public ForwardClientFactory(TransactionStateChecker transactionStateChecker) { this.init(); - this.mqClientAPIConstructor = new MQClientAPIConstructor(this.rpcHook); - this.transactionClientConstructor = new TransactionClientConstructor(this.rpcHook); + this.mqClientFactory = new MQClientFactoryImpl(this.rpcHook); + this.transactionalProducerFactory = new TransactionalProducerFactory(this.rpcHook, transactionStateChecker); } private void init() { System.setProperty(ClientConfig.SEND_MESSAGE_WITH_VIP_CHANNEL_PROPERTY, System.getProperty(ClientConfig.SEND_MESSAGE_WITH_VIP_CHANNEL_PROPERTY, "false")); - if (StringUtils.isEmpty(ConfigurationManager.getProxyConfig().getNameSrvAddr())) { + if (StringUtils.isEmpty(ConfigurationManager.getProxyConfig().getNameSrvDomain())) { System.setProperty(MixAll.NAMESRV_ADDR_PROPERTY, ConfigurationManager.getProxyConfig().getNameSrvAddr()); } else { System.setProperty("rocketmq.namesrv.domain", ConfigurationManager.getProxyConfig().getNameSrvDomain()); @@ -53,11 +51,11 @@ public class ClientFactory implements StartAndShutdown { } public MQClientAPIExtImpl getMQClient(String instanceName, int bootstrapWorkerThreads) { - return mqClientAPIConstructor.getOne(instanceName, bootstrapWorkerThreads); + return mqClientFactory.getOne(instanceName, bootstrapWorkerThreads); } - public MQClientAPIExtImpl getTransactionClient(String instanceName, int bootstrapWorkerThreads) { - return transactionClientConstructor.getOne(instanceName, bootstrapWorkerThreads); + public MQClientAPIExtImpl getTransactionalProducer(String instanceName, int bootstrapWorkerThreads) { + return transactionalProducerFactory.getOne(instanceName, bootstrapWorkerThreads); } public void setRpcHook(RPCHook rpcHook) { @@ -71,7 +69,7 @@ public class ClientFactory implements StartAndShutdown { @Override public void shutdown() throws Exception { - this.mqClientAPIConstructor.shutdownAll(); - this.transactionClientConstructor.shutdownAll(); + this.mqClientFactory.shutdownAll(); + this.transactionalProducerFactory.shutdownAll(); } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/mqconstructor/RocketMQClientConstructor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/factory/MQClientFactory.java similarity index 89% rename from proxy/src/main/java/org/apache/rocketmq/proxy/client/mqconstructor/RocketMQClientConstructor.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/client/factory/MQClientFactory.java index 2dd0378e7e..de43082292 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/mqconstructor/RocketMQClientConstructor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/factory/MQClientFactory.java @@ -14,11 +14,9 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.rocketmq.proxy.client.mqconstructor; - -public interface RocketMQClientConstructor { +package org.apache.rocketmq.proxy.client.factory; +public interface MQClientFactory { T getOne(String instanceName, int bootstrapWorkerThreads); - void shutdownAll(); } \ No newline at end of file diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/mqconstructor/MQClientAPIConstructor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/factory/MQClientFactoryImpl.java similarity index 89% rename from proxy/src/main/java/org/apache/rocketmq/proxy/client/mqconstructor/MQClientAPIConstructor.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/client/factory/MQClientFactoryImpl.java index f3e276c7a7..d4a4979c24 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/mqconstructor/MQClientAPIConstructor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/factory/MQClientFactoryImpl.java @@ -14,16 +14,16 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.rocketmq.proxy.client.mqconstructor; +package org.apache.rocketmq.proxy.client.factory; import org.apache.rocketmq.client.ClientConfig; import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; import org.apache.rocketmq.proxy.client.processor.DoNothingClientRemotingProcessor; import org.apache.rocketmq.remoting.RPCHook; -public class MQClientAPIConstructor extends AbstractRocketMQClientConstructor { +public class MQClientFactoryImpl extends AbstractMQClientFactory { - public MQClientAPIConstructor(RPCHook rpcHook) { + public MQClientFactoryImpl(RPCHook rpcHook) { super(rpcHook); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/mqconstructor/TransactionClientConstructor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/factory/TransactionalProducerFactory.java similarity index 74% rename from proxy/src/main/java/org/apache/rocketmq/proxy/client/mqconstructor/TransactionClientConstructor.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/client/factory/TransactionalProducerFactory.java index 67f8fb62bd..6168e9e358 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/mqconstructor/TransactionClientConstructor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/factory/TransactionalProducerFactory.java @@ -14,24 +14,27 @@ * See the License for the specific language governing permissions and * limitations under the License. */ -package org.apache.rocketmq.proxy.client.mqconstructor; +package org.apache.rocketmq.proxy.client.factory; import org.apache.rocketmq.client.ClientConfig; import org.apache.rocketmq.client.impl.MQClientAPIExtImpl; import org.apache.rocketmq.proxy.client.processor.ProxyClientRemotingProcessor; +import org.apache.rocketmq.proxy.client.transaction.TransactionStateChecker; import org.apache.rocketmq.remoting.RPCHook; -public class TransactionClientConstructor extends AbstractRocketMQClientConstructor { +public class TransactionalProducerFactory extends AbstractMQClientFactory { + private final TransactionStateChecker transactionStateChecker; - public TransactionClientConstructor(RPCHook rpcHook) { + public TransactionalProducerFactory(RPCHook rpcHook, TransactionStateChecker transactionStateChecker) { super(rpcHook); + this.transactionStateChecker = transactionStateChecker; } @Override MQClientAPIExtImpl newOne(String instanceName, RPCHook rpcHook, int bootstrapWorkerThreads) { return new MQClientAPIExtImpl( createNettyClientConfig(bootstrapWorkerThreads), - new ProxyClientRemotingProcessor(null), + new ProxyClientRemotingProcessor(this.transactionStateChecker), rpcHook, new ClientConfig()); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/processor/DoNothingClientRemotingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/processor/DoNothingClientRemotingProcessor.java index 2f9d100f72..74c5874370 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/processor/DoNothingClientRemotingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/processor/DoNothingClientRemotingProcessor.java @@ -23,8 +23,7 @@ import org.apache.rocketmq.remoting.protocol.RemotingCommand; public class DoNothingClientRemotingProcessor extends ClientRemotingProcessor { - public DoNothingClientRemotingProcessor( - MQClientInstance mqClientFactory) { + public DoNothingClientRemotingProcessor(MQClientInstance mqClientFactory) { super(mqClientFactory); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/processor/ProxyClientRemotingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/processor/ProxyClientRemotingProcessor.java index 07361919b0..bb5fe37ae5 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/processor/ProxyClientRemotingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/processor/ProxyClientRemotingProcessor.java @@ -31,17 +31,18 @@ import org.apache.rocketmq.remoting.exception.RemotingCommandException; import org.apache.rocketmq.remoting.protocol.RemotingCommand; public class ProxyClientRemotingProcessor extends ClientRemotingProcessor { - private final TransactionStateChecker transactionStateChecker; - public ProxyClientRemotingProcessor( - TransactionStateChecker transactionStateChecker) { + public ProxyClientRemotingProcessor(TransactionStateChecker transactionStateChecker) { super(null); this.transactionStateChecker = transactionStateChecker; } @Override - public RemotingCommand processRequest(ChannelHandlerContext ctx, RemotingCommand request) throws RemotingCommandException { + public RemotingCommand processRequest( + ChannelHandlerContext ctx, + RemotingCommand request + ) throws RemotingCommandException { if (request.getCode() == RequestCode.CHECK_TRANSACTION_STATE) { return this.checkTransactionState(ctx, request); } @@ -49,9 +50,12 @@ public class ProxyClientRemotingProcessor extends ClientRemotingProcessor { } @Override - public RemotingCommand checkTransactionState(ChannelHandlerContext ctx, - RemotingCommand request) throws RemotingCommandException { - final CheckTransactionStateRequestHeader requestHeader = (CheckTransactionStateRequestHeader) request.decodeCommandCustomHeader(CheckTransactionStateRequestHeader.class); + public RemotingCommand checkTransactionState( + ChannelHandlerContext ctx, + RemotingCommand request + ) throws RemotingCommandException { + final CheckTransactionStateRequestHeader requestHeader = + (CheckTransactionStateRequestHeader) request.decodeCommandCustomHeader(CheckTransactionStateRequestHeader.class); final ByteBuffer byteBuffer = ByteBuffer.wrap(request.getBody()); final MessageExt messageExt = MessageDecoder.decode(byteBuffer, true, false, false); if (messageExt != null) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/AddressableMessageQueue.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/AddressableMessageQueue.java deleted file mode 100644 index ee6db196e3..0000000000 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/AddressableMessageQueue.java +++ /dev/null @@ -1,80 +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.client.route; - -import java.util.Objects; -import org.apache.rocketmq.common.message.MessageQueue; - -public class AddressableMessageQueue implements Comparable { - - private final MessageQueue messageQueue; - private final String brokerAddr; - - public AddressableMessageQueue(MessageQueue messageQueue, String brokerAddr) { - this.messageQueue = messageQueue; - this.brokerAddr = brokerAddr; - } - - @Override - public int compareTo(AddressableMessageQueue o) { - return messageQueue.compareTo(o.messageQueue); - } - - @Override - public boolean equals(Object o) { - if (this == o) { - return true; - } - if (!(o instanceof AddressableMessageQueue)) { - return false; - } - AddressableMessageQueue queue = (AddressableMessageQueue) o; - return Objects.equals(messageQueue, queue.messageQueue); - } - - @Override - public int hashCode() { - return messageQueue == null ? 1 : messageQueue.hashCode(); - } - - public int getQueueId() { - return this.messageQueue.getQueueId(); - } - - public String getBrokerName() { - return this.messageQueue.getBrokerName(); - } - - public String getTopic() { - return messageQueue.getTopic(); - } - - public MessageQueue getMessageQueue() { - return messageQueue; - } - - public String getBrokerAddr() { - return brokerAddr; - } - - @Override public String toString() { - return "AddressableMessageQueue{" + - "messageQueue=" + messageQueue + - ", brokerAddr='" + brokerAddr + '\'' + - '}'; - } -} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/MessageQueueSelector.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/MessageQueueSelector.java new file mode 100644 index 0000000000..d88f560a31 --- /dev/null +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/MessageQueueSelector.java @@ -0,0 +1,229 @@ +/* + * 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.client.route; + +import java.util.ArrayList; +import java.util.Collections; +import java.util.HashSet; +import java.util.List; +import java.util.Map; +import java.util.Objects; +import java.util.Random; +import java.util.Set; +import java.util.concurrent.ConcurrentHashMap; +import java.util.concurrent.atomic.AtomicInteger; +import java.util.stream.Collectors; +import org.apache.commons.lang3.StringUtils; +import org.apache.rocketmq.common.constant.PermName; +import org.apache.rocketmq.common.message.MessageQueue; +import org.apache.rocketmq.common.protocol.route.QueueData; + +public class MessageQueueSelector { + private static final int BROKER_ACTING_QUEUE_ID = -1; + + // multiple queues for one broker, with queueId : normal + private final List queues = new ArrayList<>(); + // one queue for one broker, with queueId : -1 + private final List brokerActingQueues = new ArrayList<>(); + private final Map brokerNameQueueMap = new ConcurrentHashMap<>(); + private final AtomicInteger queueIndex; + private final AtomicInteger brokerIndex; + + public MessageQueueSelector(TopicRouteWrapper topicRouteWrapper, boolean read) { + if (read) { + this.queues.addAll(buildRead(topicRouteWrapper)); + } else { + this.queues.addAll(buildWrite(topicRouteWrapper)); + } + buildBrokerActingQueues(topicRouteWrapper.getTopicName(), this.queues); + + this.queueIndex = new AtomicInteger(Math.abs(new Random().nextInt())); + this.brokerIndex = new AtomicInteger(Math.abs(new Random().nextInt())); + } + + private static List buildRead(TopicRouteWrapper topicRoute) { + Set queueSet = new HashSet<>(); + List qds = topicRoute.getQueueDatas(); + if (qds == null) { + return new ArrayList<>(); + } + + for (QueueData qd : qds) { + if (PermName.isReadable(qd.getPerm())) { + String brokerAddr = topicRoute.getMasterAddrPrefer(qd.getBrokerName()); + if (brokerAddr == null) { + continue; + } + + for (int i = 0; i < qd.getReadQueueNums(); i++) { + SelectableMessageQueue mq = new SelectableMessageQueue( + new MessageQueue(topicRoute.getTopicName(), qd.getBrokerName(), i), + brokerAddr); + queueSet.add(mq); + } + } + } + + return queueSet.stream().sorted().collect(Collectors.toList()); + } + + private static List buildWrite(TopicRouteWrapper topicRoute) { + Set queueSet = new HashSet<>(); + // order topic route. + if (StringUtils.isNotBlank(topicRoute.getOrderTopicConf())) { + String[] brokers = topicRoute.getOrderTopicConf().split(";"); + for (String broker : brokers) { + String[] item = broker.split(":"); + String brokerName = item[0]; + String brokerAddr = topicRoute.getMasterAddr(brokerName); + if (brokerAddr == null) { + continue; + } + + int nums = Integer.parseInt(item[1]); + for (int i = 0; i < nums; i++) { + SelectableMessageQueue mq = new SelectableMessageQueue( + new MessageQueue(topicRoute.getTopicName(), brokerName, i), + brokerAddr); + queueSet.add(mq); + } + } + } else { + List qds = topicRoute.getQueueDatas(); + if (qds == null) { + return new ArrayList<>(); + } + + for (QueueData qd : qds) { + if (PermName.isWriteable(qd.getPerm())) { + String brokerAddr = topicRoute.getMasterAddr(qd.getBrokerName()); + if (brokerAddr == null) { + continue; + } + + for (int i = 0; i < qd.getWriteQueueNums(); i++) { + SelectableMessageQueue mq = new SelectableMessageQueue( + new MessageQueue(topicRoute.getTopicName(), qd.getBrokerName(), i), + brokerAddr); + queueSet.add(mq); + } + } + } + } + + return queueSet.stream().sorted().collect(Collectors.toList()); + } + + private void buildBrokerActingQueues(String topic, List normalQueues) { + for (SelectableMessageQueue mq : normalQueues) { + SelectableMessageQueue brokerActingQueue = new SelectableMessageQueue( + new MessageQueue(topic, mq.getMessageQueue().getBrokerName(), BROKER_ACTING_QUEUE_ID), + mq.getBrokerAddr()); + + if (!brokerActingQueues.contains(brokerActingQueue)) { + brokerActingQueues.add(brokerActingQueue); + brokerNameQueueMap.put(brokerActingQueue.getBrokerName(), brokerActingQueue); + } + } + + Collections.sort(brokerActingQueues); + } + + public final SelectableMessageQueue getQueueByBrokerName(String brokerName) { + return this.brokerNameQueueMap.get(brokerName); + } + + public final SelectableMessageQueue selectOne(boolean onlyBroker) { + int nextIndex = onlyBroker ? brokerIndex.getAndIncrement() : queueIndex.getAndIncrement(); + return selectOneByIndex(nextIndex, onlyBroker); + } + + public final SelectableMessageQueue selectOne(String brokerName, int queueId) { + for (SelectableMessageQueue addressableMessageQueue : queues) { + String queueBrokerName = addressableMessageQueue.getBrokerName(); + if (queueBrokerName.equals(brokerName) && addressableMessageQueue.getQueueId() == queueId) { + return addressableMessageQueue; + } + } + return null; + } + + public final SelectableMessageQueue selectOneByIndex(int index, boolean onlyBroker) { + if (onlyBroker) { + if (brokerActingQueues.isEmpty()) { + return null; + } + return brokerActingQueues.get(Math.abs(index) % brokerActingQueues.size()); + } + + if (queues.isEmpty()) { + return null; + } + return queues.get(Math.abs(index) % queues.size()); + } + + // find next same type(but different) queue with last(normal queue or broker acting queue). + public final SelectableMessageQueue selectNextQueue(SelectableMessageQueue last) { + boolean onlyBroker = last.getQueueId() < 0; + SelectableMessageQueue newOne = last; + int count = onlyBroker ? brokerActingQueues.size() : queues.size(); + + for (int i = 0; i < count; i++) { + newOne = selectOne(onlyBroker); + if (!newOne.getBrokerName().equals(last.getBrokerName()) || newOne.getQueueId() != last.getQueueId()) { + break; + } + } + + return newOne; + } + + public List getQueues() { + return queues; + } + + public List getBrokerActingQueues() { + return brokerActingQueues; + } + + @Override + public boolean equals(Object o) { + if (this == o) { + return true; + } + if (!(o instanceof MessageQueueSelector)) { + return false; + } + MessageQueueSelector queue = (MessageQueueSelector) o; + return Objects.equals(queues, queue.queues) && + Objects.equals(brokerActingQueues, queue.brokerActingQueues); + } + + @Override + public int hashCode() { + return Objects.hash(queues, brokerActingQueues); + } + + @Override + public String toString() { + return "SelectableMessageQueue{" + "queues=" + queues + + ", brokers=" + brokerActingQueues + + ", queueIndex=" + queueIndex + + ", brokerIndex=" + brokerIndex + + '}'; + } +} diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/MessageQueueWrapper.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/MessageQueueWrapper.java index d93eaab979..50bb80959d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/MessageQueueWrapper.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/MessageQueueWrapper.java @@ -19,17 +19,17 @@ package org.apache.rocketmq.proxy.client.route; import org.apache.rocketmq.common.protocol.route.TopicRouteData; public class MessageQueueWrapper { - public static final MessageQueueWrapper EMPTY_CACHED_QUEUE = new MessageQueueWrapper("", new TopicRouteData()); + public static final MessageQueueWrapper WRAPPED_EMPTY_QUEUE = new MessageQueueWrapper("", new TopicRouteData()); - private final SelectableMessageQueue read; - private final SelectableMessageQueue write; + private final MessageQueueSelector readSelector; + private final MessageQueueSelector writeSelector; private final TopicRouteWrapper topicRouteWrapper; public MessageQueueWrapper(String topic, TopicRouteData topicRouteData) { this.topicRouteWrapper = new TopicRouteWrapper(topicRouteData, topic); - this.read = new SelectableMessageQueue(topicRouteWrapper, true); - this.write = new SelectableMessageQueue(topicRouteWrapper, false); + this.readSelector = new MessageQueueSelector(topicRouteWrapper, true); + this.writeSelector = new MessageQueueSelector(topicRouteWrapper, false); } public TopicRouteData getTopicRouteData() { @@ -41,22 +41,22 @@ public class MessageQueueWrapper { } public boolean isEmptyCachedQueue() { - return this == EMPTY_CACHED_QUEUE; + return this == WRAPPED_EMPTY_QUEUE; } - public SelectableMessageQueue getRead() { - return read; + public MessageQueueSelector getReadSelector() { + return readSelector; } - public SelectableMessageQueue getWrite() { - return write; + public MessageQueueSelector getWriteSelector() { + return writeSelector; } @Override public String toString() { return "MessageQueueWrapper{" + - "read=" + read + - ", write=" + write + + "readSelector=" + readSelector + + ", writeSelector=" + writeSelector + ", topicRouteWrapper=" + topicRouteWrapper + '}'; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/SelectableMessageQueue.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/SelectableMessageQueue.java index 88f7b3cee8..7f3f47adfa 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/SelectableMessageQueue.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/route/SelectableMessageQueue.java @@ -16,188 +16,22 @@ */ package org.apache.rocketmq.proxy.client.route; -import java.util.ArrayList; -import java.util.Collections; -import java.util.HashSet; -import java.util.List; -import java.util.Map; import java.util.Objects; -import java.util.Random; -import java.util.Set; -import java.util.concurrent.ConcurrentHashMap; -import java.util.concurrent.atomic.AtomicInteger; -import java.util.stream.Collectors; -import org.apache.commons.lang3.StringUtils; -import org.apache.rocketmq.common.constant.PermName; import org.apache.rocketmq.common.message.MessageQueue; -import org.apache.rocketmq.common.protocol.route.QueueData; -public class SelectableMessageQueue { - private static final int BROKER_ACTING_QUEUE_ID = -1; +public class SelectableMessageQueue implements Comparable { - // multiple queues for one broker, with queueId : normal - private final List queues = new ArrayList<>(); - // one queue for one broker, with queueId : -1 - private final List brokerActingQueues = new ArrayList<>(); - private final Map brokerNameQueueMap = new ConcurrentHashMap<>(); - private final AtomicInteger queueIndex; - private final AtomicInteger brokerIndex; + private final MessageQueue messageQueue; + private final String brokerAddr; - public SelectableMessageQueue(TopicRouteWrapper topicRouteWrapper, boolean read) { - if (read) { - this.queues.addAll(buildRead(topicRouteWrapper)); - } else { - this.queues.addAll(buildWrite(topicRouteWrapper)); - } - buildBrokerActingQueues(topicRouteWrapper.getTopicName(), this.queues); - - this.queueIndex = new AtomicInteger(Math.abs(new Random().nextInt())); - this.brokerIndex = new AtomicInteger(Math.abs(new Random().nextInt())); + public SelectableMessageQueue(MessageQueue messageQueue, String brokerAddr) { + this.messageQueue = messageQueue; + this.brokerAddr = brokerAddr; } - private static List buildRead(TopicRouteWrapper topicRoute) { - Set queueSet = new HashSet<>(); - List qds = topicRoute.getQueueDatas(); - if (qds == null) { - return new ArrayList<>(); - } - - for (QueueData qd : qds) { - if (PermName.isReadable(qd.getPerm())) { - String brokerAddr = topicRoute.getMasterAddrPrefer(qd.getBrokerName()); - if (brokerAddr == null) { - continue; - } - - for (int i = 0; i < qd.getReadQueueNums(); i++) { - AddressableMessageQueue mq = new AddressableMessageQueue( - new MessageQueue(topicRoute.getTopicName(), qd.getBrokerName(), i), - brokerAddr); - queueSet.add(mq); - } - } - } - - return queueSet.stream().sorted().collect(Collectors.toList()); - } - - private static List buildWrite(TopicRouteWrapper topicRoute) { - Set queueSet = new HashSet<>(); - // order topic route. - if (StringUtils.isNotBlank(topicRoute.getOrderTopicConf())) { - String[] brokers = topicRoute.getOrderTopicConf().split(";"); - for (String broker : brokers) { - String[] item = broker.split(":"); - String brokerName = item[0]; - String brokerAddr = topicRoute.getMasterAddr(brokerName); - if (brokerAddr == null) { - continue; - } - - int nums = Integer.parseInt(item[1]); - for (int i = 0; i < nums; i++) { - AddressableMessageQueue mq = new AddressableMessageQueue( - new MessageQueue(topicRoute.getTopicName(), brokerName, i), - brokerAddr); - queueSet.add(mq); - } - } - } else { - List qds = topicRoute.getQueueDatas(); - if (qds == null) { - return new ArrayList<>(); - } - - for (QueueData qd : qds) { - if (PermName.isWriteable(qd.getPerm())) { - String brokerAddr = topicRoute.getMasterAddr(qd.getBrokerName()); - if (brokerAddr == null) { - continue; - } - - for (int i = 0; i < qd.getWriteQueueNums(); i++) { - AddressableMessageQueue mq = new AddressableMessageQueue( - new MessageQueue(topicRoute.getTopicName(), qd.getBrokerName(), i), - brokerAddr); - queueSet.add(mq); - } - } - } - } - - return queueSet.stream().sorted().collect(Collectors.toList()); - } - - private void buildBrokerActingQueues(String topic, List normalQueues) { - for (AddressableMessageQueue mq : normalQueues) { - AddressableMessageQueue brokerActingQueue = new AddressableMessageQueue( - new MessageQueue(topic, mq.getMessageQueue().getBrokerName(), BROKER_ACTING_QUEUE_ID), - mq.getBrokerAddr()); - - if (!brokerActingQueues.contains(brokerActingQueue)) { - brokerActingQueues.add(brokerActingQueue); - brokerNameQueueMap.put(brokerActingQueue.getBrokerName(), brokerActingQueue); - } - } - - Collections.sort(brokerActingQueues); - } - - public final AddressableMessageQueue getQueueByBrokerName(String brokerName) { - return this.brokerNameQueueMap.get(brokerName); - } - - public final AddressableMessageQueue selectOne(boolean onlyBroker) { - int nextIndex = onlyBroker ? brokerIndex.getAndIncrement() : queueIndex.getAndIncrement(); - return selectOneByIndex(nextIndex, onlyBroker); - } - - public final AddressableMessageQueue selectOne(String brokerName, int queueId) { - for (AddressableMessageQueue addressableMessageQueue : queues) { - String queueBrokerName = addressableMessageQueue.getBrokerName(); - if (queueBrokerName.equals(brokerName) && addressableMessageQueue.getQueueId() == queueId) { - return addressableMessageQueue; - } - } - return null; - } - - public final AddressableMessageQueue selectOneByIndex(int index, boolean onlyBroker) { - if (onlyBroker) { - if (brokerActingQueues.isEmpty()) { - return null; - } - return brokerActingQueues.get(Math.abs(index) % brokerActingQueues.size()); - } - - if (queues.isEmpty()) { - return null; - } - return queues.get(Math.abs(index) % queues.size()); - } - - // find next same type(but different) queue with last(normal queue or broker acting queue). - public final AddressableMessageQueue selectNextQueue(AddressableMessageQueue last) { - boolean onlyBroker = last.getQueueId() < 0; - AddressableMessageQueue newOne = last; - int count = onlyBroker ? brokerActingQueues.size() : queues.size(); - - for (int i = 0; i < count; i++) { - newOne = selectOne(onlyBroker); - if (!newOne.getBrokerName().equals(last.getBrokerName()) || newOne.getQueueId() != last.getQueueId()) { - break; - } - } - - return newOne; - } - - public List getQueues() { - return queues; - } - - public List getBrokerActingQueues() { - return brokerActingQueues; + @Override + public int compareTo(SelectableMessageQueue o) { + return messageQueue.compareTo(o.messageQueue); } @Override @@ -209,21 +43,38 @@ public class SelectableMessageQueue { return false; } SelectableMessageQueue queue = (SelectableMessageQueue) o; - return Objects.equals(queues, queue.queues) && - Objects.equals(brokerActingQueues, queue.brokerActingQueues); + return Objects.equals(messageQueue, queue.messageQueue); } @Override public int hashCode() { - return Objects.hash(queues, brokerActingQueues); + return messageQueue == null ? 1 : messageQueue.hashCode(); } - @Override - public String toString() { - return "SelectableMessageQueue{" + "queues=" + queues + - ", brokers=" + brokerActingQueues + - ", queueIndex=" + queueIndex + - ", brokerIndex=" + brokerIndex + + public int getQueueId() { + return this.messageQueue.getQueueId(); + } + + public String getBrokerName() { + return this.messageQueue.getBrokerName(); + } + + public String getTopic() { + return messageQueue.getTopic(); + } + + public MessageQueue getMessageQueue() { + return messageQueue; + } + + public String getBrokerAddr() { + return brokerAddr; + } + + @Override public String toString() { + return "AddressableMessageQueue{" + + "messageQueue=" + messageQueue + + ", brokerAddr='" + brokerAddr + '\'' + '}'; } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/transaction/TransactionId.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/transaction/TransactionId.java index 3c5b06de72..05da6904cd 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/transaction/TransactionId.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/transaction/TransactionId.java @@ -176,8 +176,15 @@ public class TransactionId { this.gatewayTransactionId = gatewayTransactionId; } + @Override public String toString() { - return "TransactionId(brokerAddr=" + this.getBrokerAddr() + ", brokerTransactionId=" + this.getBrokerTransactionId() + ", commitLogOffset=" + this.getCommitLogOffset() + ", tranStateTableOffset=" + this.getTranStateTableOffset() + ", gatewayTransactionId=" + this.getGatewayTransactionId() + ")"; + return "TransactionId{" + + "brokerAddr=" + brokerAddr + + ", brokerTransactionId='" + brokerTransactionId + '\'' + + ", commitLogOffset=" + commitLogOffset + + ", tranStateTableOffset=" + tranStateTableOffset + + ", gatewayTransactionId='" + gatewayTransactionId + '\'' + + '}'; } public static class TransactionIdBuilder { @@ -219,8 +226,15 @@ public class TransactionId { return new TransactionId(brokerAddr, brokerTransactionId, commitLogOffset, tranStateTableOffset, gatewayTransactionId); } + @Override public String toString() { - return "TransactionId.TransactionIdBuilder(brokerAddr=" + this.brokerAddr + ", brokerTransactionId=" + this.brokerTransactionId + ", commitLogOffset=" + this.commitLogOffset + ", tranStateTableOffset=" + this.tranStateTableOffset + ", gatewayTransactionId=" + this.gatewayTransactionId + ")"; + return "TransactionId.TransactionIdBuilder{" + + "brokerAddr=" + brokerAddr + + ", brokerTransactionId='" + brokerTransactionId + '\'' + + ", commitLogOffset=" + commitLogOffset + + ", tranStateTableOffset=" + tranStateTableOffset + + ", gatewayTransactionId='" + gatewayTransactionId + '\'' + + '}'; } } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/transaction/TransactionStateCheckRequest.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/transaction/TransactionStateCheckRequest.java index 246ddde201..73cab80586 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/transaction/TransactionStateCheckRequest.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/transaction/TransactionStateCheckRequest.java @@ -26,8 +26,14 @@ public class TransactionStateCheckRequest { private TransactionId transactionId; private MessageExt messageExt; - public TransactionStateCheckRequest(String groupId, Long tranStateTableOffset, Long commitLogOffset, - String msgId, TransactionId transactionId, MessageExt messageExt) { + public TransactionStateCheckRequest( + String groupId, + Long tranStateTableOffset, + Long commitLogOffset, + String msgId, + TransactionId transactionId, + MessageExt messageExt + ) { this.groupId = groupId; this.tranStateTableOffset = tranStateTableOffset; this.commitLogOffset = commitLogOffset; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/client/transaction/TransactionStateChecker.java b/proxy/src/main/java/org/apache/rocketmq/proxy/client/transaction/TransactionStateChecker.java index bfbb0e7c43..6cea826cb2 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/client/transaction/TransactionStateChecker.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/client/transaction/TransactionStateChecker.java @@ -17,6 +17,5 @@ package org.apache.rocketmq.proxy.client.transaction; public interface TransactionStateChecker { - void checkTransactionState(TransactionStateCheckRequest checkData); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/configuration/ProxyConfig.java b/proxy/src/main/java/org/apache/rocketmq/proxy/configuration/ProxyConfig.java index e69d31ae4c..56337eac8f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/configuration/ProxyConfig.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/configuration/ProxyConfig.java @@ -32,7 +32,7 @@ public class ProxyConfig { * configuration for ThreadPoolMonitor */ private boolean enablePrintJstack = true; - private long printJstackPeriodMillis = 60000; + private long printJstackInMillis = 60000; private String nameSrvAddr = "11.165.223.199:9876"; private String nameSrvDomain = ""; @@ -56,16 +56,16 @@ public class ProxyConfig { */ private int grpcMaxInboundMessageSize = 130 * 1024 * 1024; - private int expiredChannelTimeSec = 120; + private int channelExpiredInSeconds = 120; - private int consumerClientNum = 2; - private double consumerClientWorkerFactor = 0.2f; - private int producerClientNum = 2; - private double producerClientWorkerFactor = 0.2f; - private int defaultClientNum = 2; - private double defaultClientWorkerFactor = 0.2f; + private int forwardConsumerNum = 2; + private double forwardConsumerWorkerFactor = 0.2f; + private int forwardProducerNum = 2; + private double forwardProducerWorkerFactor = 0.2f; + private int defaultForwardClientNum = 2; + private double defaultForwardClientWorkerFactor = 0.2f; - private int topicRouteCacheExpireSecond = 20; + private int topicRouteCacheExpiredInSeconds = 20; private int topicRouteCacheExecutorThreadNum = 3; private int topicRouteCacheExecutorQueueCapacity = 1000; private int topicRouteCacheMaxNum = 20000; @@ -96,12 +96,12 @@ public class ProxyConfig { this.enablePrintJstack = enablePrintJstack; } - public long getPrintJstackPeriodMillis() { - return printJstackPeriodMillis; + public long getPrintJstackInMillis() { + return printJstackInMillis; } - public void setPrintJstackPeriodMillis(long printJstackPeriodMillis) { - this.printJstackPeriodMillis = printJstackPeriodMillis; + public void setPrintJstackInMillis(long printJstackInMillis) { + this.printJstackInMillis = printJstackInMillis; } public String getNameSrvAddr() { @@ -216,68 +216,68 @@ public class ProxyConfig { this.grpcMaxInboundMessageSize = grpcMaxInboundMessageSize; } - public int getExpiredChannelTimeSec() { - return expiredChannelTimeSec; + public int getChannelExpiredInSeconds() { + return channelExpiredInSeconds; } - public void setExpiredChannelTimeSec(int expiredChannelTimeSec) { - this.expiredChannelTimeSec = expiredChannelTimeSec; + public void setChannelExpiredInSeconds(int channelExpiredInSeconds) { + this.channelExpiredInSeconds = channelExpiredInSeconds; } - public int getConsumerClientNum() { - return consumerClientNum; + public int getForwardConsumerNum() { + return forwardConsumerNum; } - public void setConsumerClientNum(int consumerClientNum) { - this.consumerClientNum = consumerClientNum; + public void setForwardConsumerNum(int forwardConsumerNum) { + this.forwardConsumerNum = forwardConsumerNum; } - public double getConsumerClientWorkerFactor() { - return consumerClientWorkerFactor; + public double getForwardConsumerWorkerFactor() { + return forwardConsumerWorkerFactor; } - public void setConsumerClientWorkerFactor(double consumerClientWorkerFactor) { - this.consumerClientWorkerFactor = consumerClientWorkerFactor; + public void setForwardConsumerWorkerFactor(double forwardConsumerWorkerFactor) { + this.forwardConsumerWorkerFactor = forwardConsumerWorkerFactor; } - public int getProducerClientNum() { - return producerClientNum; + public int getForwardProducerNum() { + return forwardProducerNum; } - public void setProducerClientNum(int producerClientNum) { - this.producerClientNum = producerClientNum; + public void setForwardProducerNum(int forwardProducerNum) { + this.forwardProducerNum = forwardProducerNum; } - public double getProducerClientWorkerFactor() { - return producerClientWorkerFactor; + public double getForwardProducerWorkerFactor() { + return forwardProducerWorkerFactor; } - public void setProducerClientWorkerFactor(double producerClientWorkerFactor) { - this.producerClientWorkerFactor = producerClientWorkerFactor; + public void setForwardProducerWorkerFactor(double forwardProducerWorkerFactor) { + this.forwardProducerWorkerFactor = forwardProducerWorkerFactor; } - public int getDefaultClientNum() { - return defaultClientNum; + public int getDefaultForwardClientNum() { + return defaultForwardClientNum; } - public void setDefaultClientNum(int defaultClientNum) { - this.defaultClientNum = defaultClientNum; + public void setDefaultForwardClientNum(int defaultForwardClientNum) { + this.defaultForwardClientNum = defaultForwardClientNum; } - public double getDefaultClientWorkerFactor() { - return defaultClientWorkerFactor; + public double getDefaultForwardClientWorkerFactor() { + return defaultForwardClientWorkerFactor; } - public void setDefaultClientWorkerFactor(double defaultClientWorkerFactor) { - this.defaultClientWorkerFactor = defaultClientWorkerFactor; + public void setDefaultForwardClientWorkerFactor(double defaultForwardClientWorkerFactor) { + this.defaultForwardClientWorkerFactor = defaultForwardClientWorkerFactor; } - public int getTopicRouteCacheExpireSecond() { - return topicRouteCacheExpireSecond; + public int getTopicRouteCacheExpiredInSeconds() { + return topicRouteCacheExpiredInSeconds; } - public void setTopicRouteCacheExpireSecond(int topicRouteCacheExpireSecond) { - this.topicRouteCacheExpireSecond = topicRouteCacheExpireSecond; + public void setTopicRouteCacheExpiredInSeconds(int topicRouteCacheExpiredInSeconds) { + this.topicRouteCacheExpiredInSeconds = topicRouteCacheExpiredInSeconds; } public int getTopicRouteCacheExecutorThreadNum() { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java index 798a918523..a7d1404eda 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java @@ -54,7 +54,7 @@ import apache.rocketmq.v1.SendMessageResponse; import io.grpc.Context; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.common.constant.LoggerName; -import org.apache.rocketmq.proxy.client.ClientManager; +import org.apache.rocketmq.proxy.client.ForwardClientManager; import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown; import org.apache.rocketmq.proxy.grpc.service.cluster.ProducerService; import org.apache.rocketmq.proxy.grpc.service.cluster.RouteService; @@ -64,12 +64,12 @@ import org.slf4j.LoggerFactory; public class ClusterGrpcService extends AbstractStartAndShutdown implements GrpcForwardService { private static final Logger LOGGER = LoggerFactory.getLogger(LoggerName.GRPC_LOGGER_NAME); - private final ClientManager clientManager; + private final ForwardClientManager clientManager; private final ProducerService producerService; private final RouteService routeService; public ClusterGrpcService() { - this.clientManager = new ClientManager(checkData -> { + this.clientManager = new ForwardClientManager(checkData -> { }); this.producerService = new ProducerService(clientManager); this.routeService = new RouteService(clientManager); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/BaseService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/BaseService.java index acb9691322..b9837ef2a1 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/BaseService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/BaseService.java @@ -16,13 +16,13 @@ */ package org.apache.rocketmq.proxy.grpc.service.cluster; -import org.apache.rocketmq.proxy.client.ClientManager; +import org.apache.rocketmq.proxy.client.ForwardClientManager; public class BaseService { - protected final ClientManager clientManager; + protected final ForwardClientManager clientManager; - public BaseService(ClientManager clientManager) { + public BaseService(ForwardClientManager clientManager) { this.clientManager = clientManager; } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java index 7aa6887443..cc898d16ab 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java @@ -20,11 +20,11 @@ import apache.rocketmq.v1.ReceiveMessageRequest; import apache.rocketmq.v1.ReceiveMessageResponse; import io.grpc.Context; import java.util.concurrent.CompletableFuture; -import org.apache.rocketmq.proxy.client.ClientManager; +import org.apache.rocketmq.proxy.client.ForwardClientManager; public class ConsumerService extends BaseService { - public ConsumerService(ClientManager clientManager) { + public ConsumerService(ForwardClientManager clientManager) { super(clientManager); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerService.java index 94c316778d..8cabdf0204 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerService.java @@ -25,8 +25,8 @@ import org.apache.rocketmq.client.producer.SendResult; import org.apache.rocketmq.common.message.Message; import org.apache.rocketmq.common.message.MessageConst; import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader; -import org.apache.rocketmq.proxy.client.ClientManager; -import org.apache.rocketmq.proxy.client.route.AddressableMessageQueue; +import org.apache.rocketmq.proxy.client.ForwardClientManager; +import org.apache.rocketmq.proxy.client.route.SelectableMessageQueue; import org.apache.rocketmq.proxy.common.utils.ProxyUtils; import org.apache.rocketmq.proxy.grpc.common.Converter; import org.apache.rocketmq.proxy.grpc.common.ProxyException; @@ -42,19 +42,19 @@ public class ProducerService extends BaseService { private volatile ProducerServiceHook producerServiceHook = null; private volatile MessageQueueSelector messageQueueSelector = new DefaultMessageQueueSelector(); - public ProducerService(ClientManager clientManager) { + public ProducerService(ForwardClientManager clientManager) { super(clientManager); } public interface MessageQueueSelector { - AddressableMessageQueue selectQueue(Context ctx, SendMessageRequest request, SendMessageRequestHeader requestHeader, + SelectableMessageQueue selectQueue(Context ctx, SendMessageRequest request, SendMessageRequestHeader requestHeader, org.apache.rocketmq.common.message.Message message); } public class DefaultMessageQueueSelector implements MessageQueueSelector { @Override - public AddressableMessageQueue selectQueue(Context ctx, SendMessageRequest request, SendMessageRequestHeader requestHeader, + public SelectableMessageQueue selectQueue(Context ctx, SendMessageRequest request, SendMessageRequestHeader requestHeader, org.apache.rocketmq.common.message.Message message) { try { String topic = requestHeader.getTopic(); @@ -64,7 +64,7 @@ public class ProducerService extends BaseService { } Integer queueId = requestHeader.getQueueId(); String shardingKey = message.getProperty(MessageConst.PROPERTY_SHARDING_KEY); - AddressableMessageQueue addressableMessageQueue; + SelectableMessageQueue addressableMessageQueue; if (!StringUtils.isBlank(brokerName) && queueId != null) { // Grpc client sendSelect situation addressableMessageQueue = selectTargetQueue(topic, brokerName, queueId); @@ -81,24 +81,24 @@ public class ProducerService extends BaseService { } } - protected AddressableMessageQueue selectNormalQueue(String topic) throws Exception { + protected SelectableMessageQueue selectNormalQueue(String topic) throws Exception { return clientManager.getTopicRouteCache().selectOneWriteQueue(topic, null); } - protected AddressableMessageQueue selectTargetQueue(String topic, String brokerName, int queueId) throws Exception { + protected SelectableMessageQueue selectTargetQueue(String topic, String brokerName, int queueId) throws Exception { return clientManager.getTopicRouteCache().selectOneWriteQueue(topic, brokerName, queueId); } - protected AddressableMessageQueue selectOrderQueue(String topic, String shardingKey) throws Exception { + protected SelectableMessageQueue selectOrderQueue(String topic, String shardingKey) throws Exception { return clientManager.getTopicRouteCache().selectOneWriteQueueByKey(topic, shardingKey, null); } } public interface ProducerServiceHook { - void beforeSend(Context ctx, AddressableMessageQueue addressableMessageQueue, Message msg, SendMessageRequestHeader requestHeader); + void beforeSend(Context ctx, SelectableMessageQueue addressableMessageQueue, Message msg, SendMessageRequestHeader requestHeader); - void afterSend(Context ctx, AddressableMessageQueue addressableMessageQueue, Message msg, SendMessageRequestHeader requestHeader, + void afterSend(Context ctx, SelectableMessageQueue addressableMessageQueue, Message msg, SendMessageRequestHeader requestHeader, SendResult sendResult); } @@ -116,7 +116,7 @@ public class ProducerService extends BaseService { try { SendMessageRequestHeader requestHeader = Converter.buildSendMessageRequestHeader(request); - AddressableMessageQueue addressableMessageQueue = messageQueueSelector.selectQueue(ctx, request, requestHeader, message); + SelectableMessageQueue addressableMessageQueue = messageQueueSelector.selectQueue(ctx, request, requestHeader, message); String topic = requestHeader.getTopic(); if (addressableMessageQueue == null) { @@ -127,7 +127,7 @@ public class ProducerService extends BaseService { if (producerServiceHook != null) { producerServiceHook.beforeSend(ctx, addressableMessageQueue, message, requestHeader); } - CompletableFuture sendResultCompletableFuture = this.clientManager.getProducerClient().sendMessage( + CompletableFuture sendResultCompletableFuture = this.clientManager.getForwardProducer().sendMessage( addressableMessageQueue.getBrokerAddr(), addressableMessageQueue.getBrokerName(), message, diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java index 6bb10534bc..3809e8073e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java @@ -34,8 +34,8 @@ import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.common.constant.PermName; import org.apache.rocketmq.common.protocol.route.QueueData; import org.apache.rocketmq.common.protocol.route.TopicRouteData; -import org.apache.rocketmq.proxy.client.ClientManager; -import org.apache.rocketmq.proxy.client.route.AddressableMessageQueue; +import org.apache.rocketmq.proxy.client.ForwardClientManager; +import org.apache.rocketmq.proxy.client.route.SelectableMessageQueue; import org.apache.rocketmq.proxy.client.route.MessageQueueWrapper; import org.apache.rocketmq.proxy.common.RocketMQHelper; import org.apache.rocketmq.proxy.grpc.common.Converter; @@ -47,7 +47,7 @@ public class RouteService extends BaseService { private volatile QueryRouteHook queryRouteHook = null; private volatile QueryAssignmentHook queryAssignmentHook = null; - public RouteService(ClientManager clientManager) { + public RouteService(ForwardClientManager clientManager) { super(clientManager); } @@ -61,16 +61,16 @@ public class RouteService extends BaseService { } public interface RouteAssignmentQueueSelector { - List getAssignment(QueryAssignmentRequest request) throws Exception; + List getAssignment(QueryAssignmentRequest request) throws Exception; } public class DefaultRouteAssignmentQueueSelector implements RouteAssignmentQueueSelector { @Override - public List getAssignment(QueryAssignmentRequest request) throws Exception { + public List getAssignment(QueryAssignmentRequest request) throws Exception { MessageQueueWrapper messageQueueWrapper = clientManager.getTopicRouteCache() .getMessageQueue(Converter.getResourceNameWithNamespace(request.getTopic())); - return messageQueueWrapper.getRead().getBrokerActingQueues(); + return messageQueueWrapper.getReadSelector().getBrokerActingQueues(); } } @@ -189,9 +189,9 @@ public class RouteService extends BaseService { } List assignments = new ArrayList<>(); - List messageQueueList = this.assignmentQueueSelector.getAssignment(request); + List messageQueueList = this.assignmentQueueSelector.getAssignment(request); - for (AddressableMessageQueue messageQueue : messageQueueList) { + for (SelectableMessageQueue messageQueue : messageQueueList) { Broker broker = Broker.newBuilder() .setName(messageQueue.getBrokerName()) .setId(0) diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/client/ClientManagerTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/client/ClientManagerTest.java index 9a5703aa36..d3121f32e5 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/client/ClientManagerTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/client/ClientManagerTest.java @@ -30,20 +30,20 @@ public class ClientManagerTest extends InitConfigurationTest { @Test public void testClientManager() throws Exception { TransactionStateChecker mockedTransactionStateChecker = Mockito.mock(TransactionStateChecker.class); - ClientManager clientManager = new ClientManager(mockedTransactionStateChecker); + ForwardClientManager clientManager = new ForwardClientManager(mockedTransactionStateChecker); clientManager.start(); - assertThat(clientManager.getDefaultClient()).isNotNull(); - assertThat(clientManager.getDefaultClient().getClientNum()) - .isEqualTo(ConfigurationManager.getProxyConfig().getDefaultClientNum()); + assertThat(clientManager.getDefaultForwardClient()).isNotNull(); + assertThat(clientManager.getDefaultForwardClient().getClientNum()) + .isEqualTo(ConfigurationManager.getProxyConfig().getDefaultForwardClientNum()); - assertThat(clientManager.getProducerClient()).isNotNull(); - assertThat(clientManager.getProducerClient().getClientNum()) - .isEqualTo(ConfigurationManager.getProxyConfig().getProducerClientNum()); + assertThat(clientManager.getForwardProducer()).isNotNull(); + assertThat(clientManager.getForwardProducer().getClientNum()) + .isEqualTo(ConfigurationManager.getProxyConfig().getForwardProducerNum()); - assertThat(clientManager.getReadConsumerClient()).isNotNull(); - assertThat(clientManager.getReadConsumerClient().getClientNum()) - .isEqualTo(ConfigurationManager.getProxyConfig().getConsumerClientNum()); + assertThat(clientManager.getForwardReadConsumer()).isNotNull(); + assertThat(clientManager.getForwardReadConsumer().getClientNum()) + .isEqualTo(ConfigurationManager.getProxyConfig().getForwardConsumerNum()); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/BaseServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/BaseServiceTest.java index fd220378a0..9a6ab623ae 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/BaseServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/BaseServiceTest.java @@ -16,12 +16,12 @@ */ package org.apache.rocketmq.proxy.grpc.service.cluster; -import org.apache.rocketmq.proxy.client.ClientManager; -import org.apache.rocketmq.proxy.client.DefaultClient; -import org.apache.rocketmq.proxy.client.ProducerClient; -import org.apache.rocketmq.proxy.client.ReadConsumerClient; +import org.apache.rocketmq.proxy.client.ForwardClientManager; +import org.apache.rocketmq.proxy.client.DefaultForwardClient; +import org.apache.rocketmq.proxy.client.ForwardProducer; +import org.apache.rocketmq.proxy.client.ForwardReadConsumer; import org.apache.rocketmq.proxy.client.TopicRouteCache; -import org.apache.rocketmq.proxy.client.WriteConsumerClient; +import org.apache.rocketmq.proxy.client.ForwardWriteConsumer; import org.junit.Before; import org.junit.Ignore; import org.junit.runner.RunWith; @@ -35,24 +35,24 @@ import static org.mockito.Mockito.when; public abstract class BaseServiceTest { @Mock - protected ClientManager clientManager; + protected ForwardClientManager clientManager; @Mock - protected DefaultClient defaultClient; + protected DefaultForwardClient defaultClient; @Mock - protected ProducerClient producerClient; + protected ForwardProducer producerClient; @Mock - protected ReadConsumerClient readConsumerClient; + protected ForwardReadConsumer readConsumerClient; @Mock - protected WriteConsumerClient writeConsumerClient; + protected ForwardWriteConsumer writeConsumerClient; @Mock protected TopicRouteCache topicRouteCache; @Before public void before() throws Throwable { - when(clientManager.getDefaultClient()).thenReturn(defaultClient); - when(clientManager.getProducerClient()).thenReturn(producerClient); - when(clientManager.getReadConsumerClient()).thenReturn(readConsumerClient); - when(clientManager.getWriteConsumerClient()).thenReturn(writeConsumerClient); + when(clientManager.getDefaultForwardClient()).thenReturn(defaultClient); + when(clientManager.getForwardProducer()).thenReturn(producerClient); + when(clientManager.getForwardReadConsumer()).thenReturn(readConsumerClient); + when(clientManager.getForwardWriteConsumer()).thenReturn(writeConsumerClient); when(clientManager.getTopicRouteCache()).thenReturn(topicRouteCache); beforeEach(); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerServiceTest.java index 4165ad8d69..2ec8bb4602 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerServiceTest.java @@ -35,7 +35,7 @@ import org.apache.rocketmq.client.producer.SendStatus; import org.apache.rocketmq.common.message.MessageConst; import org.apache.rocketmq.common.message.MessageQueue; import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader; -import org.apache.rocketmq.proxy.client.route.AddressableMessageQueue; +import org.apache.rocketmq.proxy.client.route.SelectableMessageQueue; import org.apache.rocketmq.proxy.grpc.common.ProxyException; import org.apache.rocketmq.proxy.grpc.common.ProxyResponseCode; import org.junit.Test; @@ -56,19 +56,19 @@ public class ProducerServiceTest extends BaseServiceTest { @Override public void beforeEach() throws Throwable { - AddressableMessageQueue queue = new AddressableMessageQueue( + SelectableMessageQueue queue = new SelectableMessageQueue( new MessageQueue("topic", "selectOrderQueue", 0), "selectOrderQueueAddr"); when(topicRouteCache.selectOneWriteQueueByKey(anyString(), anyString(), isNull())) .thenReturn(queue); - queue = new AddressableMessageQueue( + queue = new SelectableMessageQueue( new MessageQueue("topic", "selectTargetQueue", 0), "selectTargetQueueAddr"); when(topicRouteCache.selectOneWriteQueue(anyString(), anyString(), anyInt())) .thenReturn(queue); - queue = new AddressableMessageQueue( + queue = new SelectableMessageQueue( new MessageQueue("topic", "selectNormalQueue", 0), "selectNormalQueueAddr"); when(topicRouteCache.selectOneWriteQueue(anyString(), isNull())) @@ -85,17 +85,17 @@ public class ProducerServiceTest extends BaseServiceTest { ProducerService producerService = new ProducerService(this.clientManager); - AtomicReference selectQueueRef = new AtomicReference<>(); + AtomicReference selectQueueRef = new AtomicReference<>(); AtomicReference messageRef = new AtomicReference<>(); producerService.setProducerServiceHook(new ProducerService.ProducerServiceHook() { @Override - public void beforeSend(Context ctx, AddressableMessageQueue addressableMessageQueue, + public void beforeSend(Context ctx, SelectableMessageQueue addressableMessageQueue, org.apache.rocketmq.common.message.Message msg, SendMessageRequestHeader requestHeader) { selectQueueRef.set(addressableMessageQueue); } @Override - public void afterSend(Context ctx, AddressableMessageQueue addressableMessageQueue, + public void afterSend(Context ctx, SelectableMessageQueue addressableMessageQueue, org.apache.rocketmq.common.message.Message msg, SendMessageRequestHeader requestHeader, SendResult sendResult) { @@ -138,17 +138,17 @@ public class ProducerServiceTest extends BaseServiceTest { ProducerService producerService = new ProducerService(this.clientManager); - AtomicReference selectQueueRef = new AtomicReference<>(); + AtomicReference selectQueueRef = new AtomicReference<>(); AtomicReference messageRef = new AtomicReference<>(); producerService.setProducerServiceHook(new ProducerService.ProducerServiceHook() { @Override - public void beforeSend(Context ctx, AddressableMessageQueue addressableMessageQueue, + public void beforeSend(Context ctx, SelectableMessageQueue addressableMessageQueue, org.apache.rocketmq.common.message.Message msg, SendMessageRequestHeader requestHeader) { selectQueueRef.set(addressableMessageQueue); } @Override - public void afterSend(Context ctx, AddressableMessageQueue addressableMessageQueue, + public void afterSend(Context ctx, SelectableMessageQueue addressableMessageQueue, org.apache.rocketmq.common.message.Message msg, SendMessageRequestHeader requestHeader, SendResult sendResult) { @@ -190,17 +190,17 @@ public class ProducerServiceTest extends BaseServiceTest { ProducerService producerService = new ProducerService(this.clientManager); - AtomicReference selectQueueRef = new AtomicReference<>(); + AtomicReference selectQueueRef = new AtomicReference<>(); AtomicReference messageRef = new AtomicReference<>(); producerService.setProducerServiceHook(new ProducerService.ProducerServiceHook() { @Override - public void beforeSend(Context ctx, AddressableMessageQueue addressableMessageQueue, + public void beforeSend(Context ctx, SelectableMessageQueue addressableMessageQueue, org.apache.rocketmq.common.message.Message msg, SendMessageRequestHeader requestHeader) { selectQueueRef.set(addressableMessageQueue); } @Override - public void afterSend(Context ctx, AddressableMessageQueue addressableMessageQueue, + public void afterSend(Context ctx, SelectableMessageQueue addressableMessageQueue, org.apache.rocketmq.common.message.Message msg, SendMessageRequestHeader requestHeader, SendResult sendResult) {