mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Do some refactoring work.
This commit is contained in:
@@ -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;
|
||||
}
|
||||
|
||||
@@ -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;
|
||||
|
||||
@@ -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() {
|
||||
|
||||
@@ -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<Long> 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<Long> 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<Long> searchOffset(
|
||||
|
||||
@@ -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<Integer> heartBeat(String heartbeatAddr, HeartbeatData heartbeatData, long timeout) throws Exception {
|
||||
return this.getClient().sendHeartbeat(heartbeatAddr, heartbeatData, timeout);
|
||||
public CompletableFuture<Integer> heartBeat(String brokerAddr, HeartbeatData heartbeatData) throws Exception {
|
||||
return this.heartBeat(brokerAddr, heartbeatData, DEFAULT_MQ_CLIENT_TIMEOUT);
|
||||
}
|
||||
public CompletableFuture<Integer> 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<SendResult> sendMessage(
|
||||
@@ -89,7 +95,7 @@ public class ForwardProducer extends AbstractForwardClient {
|
||||
}
|
||||
|
||||
public CompletableFuture<RemotingCommand> 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<RemotingCommand> sendMessageBack(String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader, long timeoutMillis) {
|
||||
|
||||
@@ -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<PopResult> 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<PopResult> popMessage(
|
||||
@@ -62,7 +61,7 @@ public class ForwardReadConsumer extends AbstractForwardClient {
|
||||
}
|
||||
|
||||
public CompletableFuture<PullResult> 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<PullResult> pullMessage(String address, PullMessageRequestHeader requestHeader,
|
||||
|
||||
@@ -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<AckResult> 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<AckResult> 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<AckResult> changeInvisibleTimeAsync(
|
||||
|
||||
+3
-3
@@ -24,14 +24,14 @@ import org.apache.rocketmq.remoting.netty.NettyClientConfig;
|
||||
import org.slf4j.Logger;
|
||||
import org.slf4j.LoggerFactory;
|
||||
|
||||
public abstract class AbstractClientFactory<T> {
|
||||
private static final Logger log = LoggerFactory.getLogger(AbstractClientFactory.class);
|
||||
public abstract class AbstractClientManager<T> {
|
||||
private static final Logger log = LoggerFactory.getLogger(AbstractClientManager.class);
|
||||
|
||||
protected final ScheduledExecutorService scheduledExecutorService;
|
||||
protected Map<String, T> cacheTable = new ConcurrentHashMap<>();
|
||||
protected RPCHook rpcHook;
|
||||
|
||||
public AbstractClientFactory(ScheduledExecutorService scheduledExecutorService, RPCHook rpcHook) {
|
||||
public AbstractClientManager(ScheduledExecutorService scheduledExecutorService, RPCHook rpcHook) {
|
||||
this.scheduledExecutorService = scheduledExecutorService;
|
||||
this.rpcHook = rpcHook;
|
||||
}
|
||||
+1
-1
@@ -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<MQClientAPIExt> {
|
||||
public abstract class AbstractMQClientFactory extends AbstractClientManager<MQClientAPIExt> {
|
||||
|
||||
public AbstractMQClientFactory(ScheduledExecutorService scheduledExecutorService, RPCHook rpcHook) {
|
||||
super(scheduledExecutorService, rpcHook);
|
||||
|
||||
+2
-2
@@ -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(
|
||||
@@ -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);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
+1
-2
@@ -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);
|
||||
}
|
||||
|
||||
+2
-3
@@ -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());
|
||||
|
||||
+1
-2
@@ -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()
|
||||
|
||||
Reference in New Issue
Block a user