[ISSUE #3949] do some renaming work.

This commit is contained in:
Jixiang.jjx
2022-07-13 11:29:09 +08:00
committed by zhouxiang
parent be512d5747
commit 3b437435b4
33 changed files with 577 additions and 548 deletions
@@ -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");
}
}
@@ -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<Map.Entry<String, SimpleChannel>> iterator = clientIdChannelMap.entrySet()
.iterator();
Iterator<Map.Entry<String, SimpleChannel>> iterator = clientIdChannelMap.entrySet().iterator();
while (iterator.hasNext()) {
Map.Entry<String, SimpleChannel> entry = iterator.next();
if (!entry.getValue()
@@ -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);
}
}
@@ -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;
}
}
@@ -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<List<String>> getConsumerListByGroup(String brokerAddr, GetConsumerListByGroupRequestHeader requestHeader,
long timeoutMillis) {
public CompletableFuture<List<String>> 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);
}
}
@@ -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;
}
}
@@ -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
@@ -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);
@@ -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);
@@ -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<String /* topicName */, MessageQueueWrapper> 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<AddressableMessageQueue> writeQueues = getMessageQueue(topic).getWrite().getQueues();
public SelectableMessageQueue selectOneWriteQueueByKey(String topic, String shardingKey, SelectableMessageQueue last) throws Exception {
List<SelectableMessageQueue> 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;
}
@@ -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<T> implements RocketMQClientConstructor<T> {
public abstract class AbstractMQClientFactory<T> implements MQClientFactory<T> {
private static final Logger log = LoggerFactory.getLogger(AbstractRocketMQClientConstructor.class);
private static final Logger LOGGER = LoggerFactory.getLogger(AbstractMQClientFactory.class);
protected Map<String, T> 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<T> implements RocketMQCl
try {
this.shutdown(v);
} catch (Exception e) {
log.warn("RocketMQClientConstructor shutdown all err.", e);
LOGGER.warn("RocketMQClientConstructor shutdown all err.", e);
}
});
}
@@ -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();
}
}
@@ -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<T> {
package org.apache.rocketmq.proxy.client.factory;
public interface MQClientFactory<T> {
T getOne(String instanceName, int bootstrapWorkerThreads);
void shutdownAll();
}
@@ -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<MQClientAPIExtImpl> {
public class MQClientFactoryImpl extends AbstractMQClientFactory<MQClientAPIExtImpl> {
public MQClientAPIConstructor(RPCHook rpcHook) {
public MQClientFactoryImpl(RPCHook rpcHook) {
super(rpcHook);
}
@@ -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<MQClientAPIExtImpl> {
public class TransactionalProducerFactory extends AbstractMQClientFactory<MQClientAPIExtImpl> {
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());
}
@@ -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);
}
@@ -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) {
@@ -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<AddressableMessageQueue> {
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 + '\'' +
'}';
}
}
@@ -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<SelectableMessageQueue> queues = new ArrayList<>();
// one queue for one broker, with queueId : -1
private final List<SelectableMessageQueue> brokerActingQueues = new ArrayList<>();
private final Map<String, SelectableMessageQueue> 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<SelectableMessageQueue> buildRead(TopicRouteWrapper topicRoute) {
Set<SelectableMessageQueue> queueSet = new HashSet<>();
List<QueueData> 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<SelectableMessageQueue> buildWrite(TopicRouteWrapper topicRoute) {
Set<SelectableMessageQueue> 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<QueueData> 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<SelectableMessageQueue> 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<SelectableMessageQueue> getQueues() {
return queues;
}
public List<SelectableMessageQueue> 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 +
'}';
}
}
@@ -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 +
'}';
}
@@ -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<SelectableMessageQueue> {
// multiple queues for one broker, with queueId : normal
private final List<AddressableMessageQueue> queues = new ArrayList<>();
// one queue for one broker, with queueId : -1
private final List<AddressableMessageQueue> brokerActingQueues = new ArrayList<>();
private final Map<String, AddressableMessageQueue> 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<AddressableMessageQueue> buildRead(TopicRouteWrapper topicRoute) {
Set<AddressableMessageQueue> queueSet = new HashSet<>();
List<QueueData> 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<AddressableMessageQueue> buildWrite(TopicRouteWrapper topicRoute) {
Set<AddressableMessageQueue> 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<QueueData> 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<AddressableMessageQueue> 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<AddressableMessageQueue> getQueues() {
return queues;
}
public List<AddressableMessageQueue> 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 + '\'' +
'}';
}
}
@@ -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 + '\'' +
'}';
}
}
}
@@ -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;
@@ -17,6 +17,5 @@
package org.apache.rocketmq.proxy.client.transaction;
public interface TransactionStateChecker {
void checkTransactionState(TransactionStateCheckRequest checkData);
}
@@ -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() {
@@ -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);
@@ -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;
}
}
@@ -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);
}
@@ -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<SendResult> sendResultCompletableFuture = this.clientManager.getProducerClient().sendMessage(
CompletableFuture<SendResult> sendResultCompletableFuture = this.clientManager.getForwardProducer().sendMessage(
addressableMessageQueue.getBrokerAddr(),
addressableMessageQueue.getBrokerName(),
message,
@@ -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<AddressableMessageQueue> getAssignment(QueryAssignmentRequest request) throws Exception;
List<SelectableMessageQueue> getAssignment(QueryAssignmentRequest request) throws Exception;
}
public class DefaultRouteAssignmentQueueSelector implements RouteAssignmentQueueSelector {
@Override
public List<AddressableMessageQueue> getAssignment(QueryAssignmentRequest request) throws Exception {
public List<SelectableMessageQueue> 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<Assignment> assignments = new ArrayList<>();
List<AddressableMessageQueue> messageQueueList = this.assignmentQueueSelector.getAssignment(request);
List<SelectableMessageQueue> messageQueueList = this.assignmentQueueSelector.getAssignment(request);
for (AddressableMessageQueue messageQueue : messageQueueList) {
for (SelectableMessageQueue messageQueue : messageQueueList) {
Broker broker = Broker.newBuilder()
.setName(messageQueue.getBrokerName())
.setId(0)
@@ -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());
@@ -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();
@@ -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<AddressableMessageQueue> selectQueueRef = new AtomicReference<>();
AtomicReference<SelectableMessageQueue> selectQueueRef = new AtomicReference<>();
AtomicReference<org.apache.rocketmq.common.message.Message> 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<AddressableMessageQueue> selectQueueRef = new AtomicReference<>();
AtomicReference<SelectableMessageQueue> selectQueueRef = new AtomicReference<>();
AtomicReference<org.apache.rocketmq.common.message.Message> 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<AddressableMessageQueue> selectQueueRef = new AtomicReference<>();
AtomicReference<SelectableMessageQueue> selectQueueRef = new AtomicReference<>();
AtomicReference<org.apache.rocketmq.common.message.Message> 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) {