diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/ProxyUtils.java b/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/ProxyUtils.java index f7c05450b6..6b9ea2cf7e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/ProxyUtils.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/common/utils/ProxyUtils.java @@ -16,10 +16,6 @@ */ package org.apache.rocketmq.proxy.common.utils; -import java.time.Duration; - public class ProxyUtils { - public static final long DEFAULT_MQ_CLIENT_TIMEOUT = Duration.ofSeconds(3).toMillis(); - public static final int MAX_MSG_NUMS_FOR_POP_REQUEST = 32; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/AbstractForwardClient.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/AbstractForwardClient.java index 1ea6539a92..c8ca3f563c 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/AbstractForwardClient.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/AbstractForwardClient.java @@ -16,25 +16,27 @@ */ package org.apache.rocketmq.proxy.connector; +import java.time.Duration; import java.util.concurrent.ThreadLocalRandom; import org.apache.rocketmq.client.impl.MQClientAPIExt; import org.apache.rocketmq.proxy.common.StartAndShutdown; -import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory; +import org.apache.rocketmq.proxy.connector.factory.ForwardClientManager; public abstract class AbstractForwardClient implements StartAndShutdown { + public static final long DEFAULT_MQ_CLIENT_TIMEOUT = Duration.ofSeconds(3).toMillis(); - private final ForwardClientFactory clientFactory; + private final ForwardClientManager clientFactory; private MQClientAPIExt[] clients; private final String gidPrefix; - public AbstractForwardClient(ForwardClientFactory clientFactory, String gidPrefix) { + public AbstractForwardClient(ForwardClientManager clientFactory, String gidPrefix) { this.clientFactory = clientFactory; this.gidPrefix = gidPrefix; } protected abstract int getClientNum(); - protected abstract MQClientAPIExt createNewClient(ForwardClientFactory forwardClientFactory, String name); + protected abstract MQClientAPIExt createNewClient(ForwardClientManager forwardClientFactory, String name); protected String getNamePrefix() { return this.gidPrefix; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ConnectorManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ConnectorManager.java index 8ae3c8d7aa..3965208d06 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ConnectorManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ConnectorManager.java @@ -16,14 +16,14 @@ */ package org.apache.rocketmq.proxy.connector; -import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory; +import org.apache.rocketmq.proxy.connector.factory.ForwardClientManager; import org.apache.rocketmq.proxy.connector.route.TopicRouteCache; import org.apache.rocketmq.proxy.connector.transaction.TransactionHeartbeatRegisterService; import org.apache.rocketmq.proxy.connector.transaction.TransactionStateChecker; import org.apache.rocketmq.proxy.common.AbstractStartAndShutdown; public class ConnectorManager extends AbstractStartAndShutdown { - private final ForwardClientFactory forwardClientFactory; + private final ForwardClientManager forwardClientManager; private final DefaultForwardClient defaultForwardClient; private final ForwardProducer forwardProducer; private final ForwardReadConsumer forwardReadConsumer; @@ -33,16 +33,16 @@ public class ConnectorManager extends AbstractStartAndShutdown { private final TransactionHeartbeatRegisterService transactionHeartbeatRegisterService; public ConnectorManager(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.forwardClientManager = new ForwardClientManager(transactionStateChecker); + this.defaultForwardClient = new DefaultForwardClient(this.forwardClientManager); + this.forwardProducer = new ForwardProducer(this.forwardClientManager); + this.forwardReadConsumer = new ForwardReadConsumer(this.forwardClientManager); + this.forwardWriteConsumer = new ForwardWriteConsumer(this.forwardClientManager); this.topicRouteCache = new TopicRouteCache(this.defaultForwardClient); this.transactionHeartbeatRegisterService = new TransactionHeartbeatRegisterService(this.forwardProducer, this.topicRouteCache); - this.appendStartAndShutdown(this.forwardClientFactory); + this.appendStartAndShutdown(this.forwardClientManager); this.appendStartAndShutdown(this.defaultForwardClient); this.appendStartAndShutdown(this.forwardProducer); this.appendStartAndShutdown(this.forwardReadConsumer); @@ -50,8 +50,8 @@ public class ConnectorManager extends AbstractStartAndShutdown { this.appendStartAndShutdown(this.transactionHeartbeatRegisterService); } - public ForwardClientFactory getForwardClientFactory() { - return forwardClientFactory; + public ForwardClientManager getForwardClientManager() { + return forwardClientManager; } public DefaultForwardClient getDefaultForwardClient() { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/DefaultForwardClient.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/DefaultForwardClient.java index 711033ae6b..cbe597e88d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/DefaultForwardClient.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/DefaultForwardClient.java @@ -22,15 +22,14 @@ import org.apache.rocketmq.client.exception.MQClientException; import org.apache.rocketmq.client.impl.MQClientAPIExt; import org.apache.rocketmq.common.protocol.header.GetConsumerListByGroupRequestHeader; import org.apache.rocketmq.common.protocol.route.TopicRouteData; -import org.apache.rocketmq.proxy.common.utils.ProxyUtils; import org.apache.rocketmq.proxy.config.ConfigurationManager; -import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory; +import org.apache.rocketmq.proxy.connector.factory.ForwardClientManager; import org.apache.rocketmq.remoting.exception.RemotingException; public class DefaultForwardClient extends AbstractForwardClient { private static final String CID_PREFIX = "CID_RMQ_PROXY_DEFAULT_"; - public DefaultForwardClient(ForwardClientFactory clientFactory) { + public DefaultForwardClient(ForwardClientManager clientFactory) { super(clientFactory, CID_PREFIX); } @@ -40,7 +39,7 @@ public class DefaultForwardClient extends AbstractForwardClient { } @Override - protected MQClientAPIExt createNewClient(ForwardClientFactory clientFactory, String name) { + protected MQClientAPIExt createNewClient(ForwardClientManager clientFactory, String name) { double workerFactor = ConfigurationManager.getProxyConfig().getDefaultForwardClientWorkerFactor(); int threadCount = (int) Math.ceil(Runtime.getRuntime().availableProcessors() * workerFactor); @@ -55,13 +54,18 @@ public class DefaultForwardClient extends AbstractForwardClient { return this.getClient().getConsumerListByGroup(brokerAddr, requestHeader, timeoutMillis); } + public TopicRouteData getTopicRouteInfoFromNameServer(String topic) + throws RemotingException, InterruptedException, MQClientException { + return this.getTopicRouteInfoFromNameServer(topic, DEFAULT_MQ_CLIENT_TIMEOUT); + } + public TopicRouteData getTopicRouteInfoFromNameServer(String topic, long timeoutMillis) throws RemotingException, InterruptedException, MQClientException { return this.getClient().getTopicRouteInfoFromNameServer(topic, timeoutMillis); } public CompletableFuture getMaxOffset(String brokerAddr, String topic, int queueId) { - return this.getMaxOffset(brokerAddr, topic, queueId, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + return this.getMaxOffset(brokerAddr, topic, queueId, DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture getMaxOffset( @@ -79,7 +83,7 @@ public class DefaultForwardClient extends AbstractForwardClient { int queueId, long timestamp ) { - return this.searchOffset(brokerAddr, topic, queueId, timestamp, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + return this.searchOffset(brokerAddr, topic, queueId, timestamp, DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture searchOffset( diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java index 3752fd6f4c..5c362713e9 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardProducer.java @@ -26,9 +26,8 @@ import org.apache.rocketmq.common.protocol.header.EndTransactionRequestHeader; import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader; import org.apache.rocketmq.common.protocol.heartbeat.HeartbeatData; import org.apache.rocketmq.common.sysflag.MessageSysFlag; -import org.apache.rocketmq.proxy.common.utils.ProxyUtils; import org.apache.rocketmq.proxy.config.ConfigurationManager; -import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory; +import org.apache.rocketmq.proxy.connector.factory.ForwardClientManager; import org.apache.rocketmq.proxy.connector.transaction.TransactionId; import org.apache.rocketmq.remoting.protocol.RemotingCommand; @@ -36,7 +35,7 @@ public class ForwardProducer extends AbstractForwardClient { private static final String PID_PREFIX = "PID_RMQ_PROXY_PUBLISH_MESSAGE_"; - public ForwardProducer(ForwardClientFactory clientFactory) { + public ForwardProducer(ForwardClientManager clientFactory) { super(clientFactory, PID_PREFIX); } @@ -46,15 +45,22 @@ public class ForwardProducer extends AbstractForwardClient { } @Override - protected MQClientAPIExt createNewClient(ForwardClientFactory clientFactory, String name) { + protected MQClientAPIExt createNewClient(ForwardClientManager clientFactory, String name) { double workerFactor = ConfigurationManager.getProxyConfig().getForwardProducerWorkerFactor(); final int threadCount = (int) Math.ceil(Runtime.getRuntime().availableProcessors() * workerFactor); return clientFactory.getTransactionalProducer(name, threadCount); } - public CompletableFuture heartBeat(String heartbeatAddr, HeartbeatData heartbeatData, long timeout) throws Exception { - return this.getClient().sendHeartbeat(heartbeatAddr, heartbeatData, timeout); + public CompletableFuture heartBeat(String brokerAddr, HeartbeatData heartbeatData) throws Exception { + return this.heartBeat(brokerAddr, heartbeatData, DEFAULT_MQ_CLIENT_TIMEOUT); + } + public CompletableFuture heartBeat(String brokerAddr, HeartbeatData heartbeatData, long timeout) throws Exception { + return this.getClient().sendHeartbeat(brokerAddr, heartbeatData, timeout); + } + + public void endTransaction(String brokerAddr, EndTransactionRequestHeader requestHeader) throws Exception { + this.endTransaction(brokerAddr, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); } public void endTransaction(String brokerAddr, EndTransactionRequestHeader requestHeader, long timeoutMillis) throws Exception { @@ -67,7 +73,7 @@ public class ForwardProducer extends AbstractForwardClient { Message msg, SendMessageRequestHeader requestHeader ) { - return this.sendMessage(address, brokerName, msg, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + return this.sendMessage(address, brokerName, msg, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture sendMessage( @@ -89,7 +95,7 @@ public class ForwardProducer extends AbstractForwardClient { } public CompletableFuture sendMessageBack(String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader) { - return this.sendMessageBack(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + return this.sendMessageBack(brokerAddr, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture sendMessageBack(String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader, long timeoutMillis) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardReadConsumer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardReadConsumer.java index 37ee10f95c..1cb9da3a6b 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardReadConsumer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardReadConsumer.java @@ -22,15 +22,14 @@ import org.apache.rocketmq.client.consumer.PullResult; import org.apache.rocketmq.client.impl.MQClientAPIExt; import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader; import org.apache.rocketmq.common.protocol.header.PullMessageRequestHeader; -import org.apache.rocketmq.proxy.common.utils.ProxyUtils; import org.apache.rocketmq.proxy.config.ConfigurationManager; -import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory; +import org.apache.rocketmq.proxy.connector.factory.ForwardClientManager; public class ForwardReadConsumer extends AbstractForwardClient { private static final String CID_PREFIX = "CID_RMQ_PROXY_CONSUME_MESSAGE_"; - public ForwardReadConsumer(ForwardClientFactory clientFactory) { + public ForwardReadConsumer(ForwardClientManager clientFactory) { super(clientFactory, CID_PREFIX); } @@ -40,7 +39,7 @@ public class ForwardReadConsumer extends AbstractForwardClient { } @Override - protected MQClientAPIExt createNewClient(ForwardClientFactory clientFactory, String name) { + protected MQClientAPIExt createNewClient(ForwardClientManager clientFactory, String name) { double workerFactor = ConfigurationManager.getProxyConfig().getForwardConsumerWorkerFactor(); final int threadCount = (int) Math.ceil(Runtime.getRuntime().availableProcessors() * workerFactor); @@ -49,7 +48,7 @@ public class ForwardReadConsumer extends AbstractForwardClient { public CompletableFuture popMessage(String address, String brokerName, PopMessageRequestHeader requestHeader) { - return this.popMessage(address, brokerName, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + return this.popMessage(address, brokerName, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture popMessage( @@ -62,7 +61,7 @@ public class ForwardReadConsumer extends AbstractForwardClient { } public CompletableFuture pullMessage(String address, PullMessageRequestHeader requestHeader) { - return this.pullMessage(address, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + return this.pullMessage(address, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture pullMessage(String address, PullMessageRequestHeader requestHeader, diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardWriteConsumer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardWriteConsumer.java index 082c8c5963..affd9a868a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardWriteConsumer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/ForwardWriteConsumer.java @@ -22,16 +22,15 @@ import org.apache.rocketmq.client.impl.MQClientAPIExt; 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.common.utils.ProxyUtils; import org.apache.rocketmq.proxy.config.ConfigurationManager; -import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory; +import org.apache.rocketmq.proxy.connector.factory.ForwardClientManager; import org.apache.rocketmq.remoting.exception.RemotingException; public class ForwardWriteConsumer extends AbstractForwardClient { private static final String CID_PREFIX = "CID_RMQ_PROXY_DELETE_MESSAGE_"; - public ForwardWriteConsumer(ForwardClientFactory clientFactory) { + public ForwardWriteConsumer(ForwardClientManager clientFactory) { super(clientFactory, CID_PREFIX); } @@ -41,7 +40,7 @@ public class ForwardWriteConsumer extends AbstractForwardClient { } @Override - protected MQClientAPIExt createNewClient(ForwardClientFactory clientFactory, String name) { + protected MQClientAPIExt createNewClient(ForwardClientManager clientFactory, String name) { double workerFactor = ConfigurationManager.getProxyConfig().getForwardConsumerWorkerFactor(); final int threadCount = (int) Math.ceil(Runtime.getRuntime().availableProcessors() * workerFactor); @@ -49,7 +48,7 @@ public class ForwardWriteConsumer extends AbstractForwardClient { } public CompletableFuture ackMessage(String address, AckMessageRequestHeader requestHeader) { - return this.ackMessage(address, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + return this.ackMessage(address, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture ackMessage( @@ -65,7 +64,7 @@ public class ForwardWriteConsumer extends AbstractForwardClient { String brokerName, ChangeInvisibleTimeRequestHeader requestHeader ) { - return this.changeInvisibleTimeAsync(address, brokerName, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + return this.changeInvisibleTimeAsync(address, brokerName, requestHeader, DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture changeInvisibleTimeAsync( diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractClientFactory.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractClientManager.java similarity index 95% rename from proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractClientFactory.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractClientManager.java index e5ea496c9a..190afd1f8a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractClientFactory.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractClientManager.java @@ -24,14 +24,14 @@ import org.apache.rocketmq.remoting.netty.NettyClientConfig; import org.slf4j.Logger; import org.slf4j.LoggerFactory; -public abstract class AbstractClientFactory { - private static final Logger log = LoggerFactory.getLogger(AbstractClientFactory.class); +public abstract class AbstractClientManager { + private static final Logger log = LoggerFactory.getLogger(AbstractClientManager.class); protected final ScheduledExecutorService scheduledExecutorService; protected Map cacheTable = new ConcurrentHashMap<>(); protected RPCHook rpcHook; - public AbstractClientFactory(ScheduledExecutorService scheduledExecutorService, RPCHook rpcHook) { + public AbstractClientManager(ScheduledExecutorService scheduledExecutorService, RPCHook rpcHook) { this.scheduledExecutorService = scheduledExecutorService; this.rpcHook = rpcHook; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractMQClientFactory.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractMQClientFactory.java index 1a772dc04d..4abcde3316 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractMQClientFactory.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/AbstractMQClientFactory.java @@ -24,7 +24,7 @@ import org.apache.rocketmq.client.impl.ClientRemotingProcessor; import org.apache.rocketmq.client.impl.MQClientAPIExt; import org.apache.rocketmq.remoting.RPCHook; -public abstract class AbstractMQClientFactory extends AbstractClientFactory { +public abstract class AbstractMQClientFactory extends AbstractClientManager { public AbstractMQClientFactory(ScheduledExecutorService scheduledExecutorService, RPCHook rpcHook) { super(scheduledExecutorService, rpcHook); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/ForwardClientFactory.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/ForwardClientManager.java similarity index 96% rename from proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/ForwardClientFactory.java rename to proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/ForwardClientManager.java index 94cd23f3b3..37ac2f2497 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/ForwardClientFactory.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/ForwardClientManager.java @@ -28,14 +28,14 @@ import org.apache.rocketmq.proxy.common.StartAndShutdown; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.remoting.RPCHook; -public class ForwardClientFactory implements StartAndShutdown { +public class ForwardClientManager implements StartAndShutdown { private RPCHook rpcHook = null; private final MQClientFactory mqClientFactory; private final TransactionProducerFactory transactionalProducerFactory; - public ForwardClientFactory(TransactionStateChecker transactionStateChecker) { + public ForwardClientManager(TransactionStateChecker transactionStateChecker) { this.init(); ScheduledExecutorService scheduledExecutorService = Executors.newSingleThreadScheduledExecutor( diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/route/TopicRouteCache.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/route/TopicRouteCache.java index b3f62c1a3d..49189a0126 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/route/TopicRouteCache.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/route/TopicRouteCache.java @@ -29,7 +29,6 @@ import org.apache.rocketmq.common.protocol.route.BrokerData; import org.apache.rocketmq.common.protocol.route.TopicRouteData; import org.apache.rocketmq.common.thread.ThreadPoolMonitor; import org.apache.rocketmq.proxy.common.AbstractCacheLoader; -import org.apache.rocketmq.proxy.common.utils.ProxyUtils; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.config.ProxyConfig; import org.apache.rocketmq.proxy.connector.DefaultForwardClient; @@ -160,7 +159,7 @@ public class TopicRouteCache { @Override protected TopicRouteData loadTopicRouteData(String topic) throws Exception { - return defaultClient.getTopicRouteInfoFromNameServer(topic, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + return defaultClient.getTopicRouteInfoFromNameServer(topic); } } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionHeartbeatRegisterService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionHeartbeatRegisterService.java index 5ae71d5f9b..d7f543b49e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionHeartbeatRegisterService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionHeartbeatRegisterService.java @@ -31,7 +31,6 @@ import org.apache.rocketmq.common.protocol.heartbeat.ProducerData; import org.apache.rocketmq.common.protocol.route.BrokerData; import org.apache.rocketmq.common.thread.ThreadPoolMonitor; import org.apache.rocketmq.proxy.common.StartAndShutdown; -import org.apache.rocketmq.proxy.common.utils.ProxyUtils; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.config.ProxyConfig; import org.apache.rocketmq.proxy.connector.ForwardProducer; @@ -156,7 +155,7 @@ public class TransactionHeartbeatRegisterService implements StartAndShutdown { heartbeatExecutors.submit(() -> { String brokerAddr = brokerData.selectBrokerAddr(); try { - this.forwardProducer.heartBeat(brokerAddr, heartbeatData, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + this.forwardProducer.heartBeat(brokerAddr, heartbeatData); } catch (Exception e) { log.error("Send transactionHeartbeat to broker err. brokerAddr: {}", brokerAddr, e); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/TransactionService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/TransactionService.java index bcc4796d8f..c55af42501 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/TransactionService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/TransactionService.java @@ -29,16 +29,15 @@ import java.util.concurrent.ThreadLocalRandom; import org.apache.commons.collections.CollectionUtils; import org.apache.rocketmq.common.protocol.header.EndTransactionRequestHeader; import org.apache.rocketmq.proxy.channel.ChannelManager; -import org.apache.rocketmq.proxy.common.utils.ProxyUtils; import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.connector.ForwardProducer; import org.apache.rocketmq.proxy.connector.transaction.TransactionId; import org.apache.rocketmq.proxy.connector.transaction.TransactionStateCheckRequest; import org.apache.rocketmq.proxy.connector.transaction.TransactionStateChecker; -import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel; import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.adapter.ResponseHook; +import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel; import org.apache.rocketmq.remoting.common.RemotingHelper; public class TransactionService extends BaseService implements TransactionStateChecker { @@ -100,7 +99,7 @@ public class TransactionService extends BaseService implements TransactionStateC TransactionId handle = TransactionId.decode(request.getTransactionId()); String brokerAddr = RemotingHelper.parseSocketAddressAddr(handle.getBrokerAddr()); EndTransactionRequestHeader requestHeader = this.toEndTransactionRequestHeader(ctx, request); - this.forwardProducer.endTransaction(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + this.forwardProducer.endTransaction(brokerAddr, requestHeader); future.complete(EndTransactionResponse.newBuilder() .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) .build()); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/TransactionServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/TransactionServiceTest.java index 86831b75d1..e54529f5e7 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/TransactionServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/cluster/TransactionServiceTest.java @@ -20,7 +20,6 @@ import org.mockito.Mock; import static org.junit.Assert.assertEquals; import static org.junit.Assert.assertTrue; import static org.mockito.ArgumentMatchers.any; -import static org.mockito.ArgumentMatchers.anyLong; import static org.mockito.ArgumentMatchers.anyString; import static org.mockito.Mockito.doAnswer; import static org.mockito.Mockito.mock; @@ -77,7 +76,7 @@ public class TransactionServiceTest extends BaseServiceTest { brokerAddrRef.set(mock.getArgument(0)); headerRef.set(mock.getArgument(1)); return null; - }).when(producerClient).endTransaction(anyString(), any(), anyLong()); + }).when(producerClient).endTransaction(anyString(), any()); EndTransactionResponse response = transactionService.endTransaction(Context.current(), EndTransactionRequest.newBuilder() .setGroup(Resource.newBuilder()