From 5288501515b15bdbf6d4e72254d9856d5c843ea0 Mon Sep 17 00:00:00 2001 From: "Jixiang.jjx" Date: Fri, 25 Mar 2022 17:59:50 +0800 Subject: [PATCH] [ISSUE #3949] Do some refactor work. --- .../apache/rocketmq/proxy/ProxyStartup.java | 6 +- .../proxy/connector/DefaultForwardClient.java | 14 ++ .../proxy/connector/ForwardProducer.java | 22 +++- .../proxy/connector/ForwardReadConsumer.java | 19 ++- .../proxy/connector/ForwardWriteConsumer.java | 21 ++- .../connector/route/MessageQueueSelector.java | 8 +- .../route/SelectableMessageQueue.java | 2 +- .../proxy/grpc/adapter/GrpcConverter.java | 13 ++ .../proxy/grpc/service/LocalGrpcService.java | 29 ++--- .../grpc/service/cluster/BaseService.java | 9 ++ .../grpc/service/cluster/ConsumerService.java | 122 +++++++++--------- .../cluster/DefaultWriteQueueSelector.java | 26 ++-- .../grpc/service/cluster/ProducerService.java | 55 ++++---- .../service/cluster/PullMessageService.java | 61 ++++----- .../grpc/service/cluster/RouteService.java | 18 ++- .../service/cluster/TransactionService.java | 14 +- 16 files changed, 243 insertions(+), 196 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java b/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java index ccacf81d84..4ed9941442 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java @@ -58,15 +58,17 @@ public class ProxyStartup { // init thread pool monitor for proxy. initThreadPoolMonitor(); - // create and start grpcServer + // create grpcServer GrpcServer grpcServer = createGrpcServer(); PROXY_START_AND_SHUTDOWN.appendStartAndShutdown(grpcServer); - // health check server + // create health check server final HealthCheckServer healthCheckServer = new HealthCheckServer(); PROXY_START_AND_SHUTDOWN.appendStartAndShutdown(healthCheckServer); + // start servers one by one. PROXY_START_AND_SHUTDOWN.start(); + Runtime.getRuntime().addShutdownHook(new Thread(() -> { LOGGER.info("try to shutdown server"); try { 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 43d7bbafd9..711033ae6b 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,6 +22,7 @@ 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.remoting.exception.RemotingException; @@ -59,6 +60,10 @@ public class DefaultForwardClient extends AbstractForwardClient { 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); + } + public CompletableFuture getMaxOffset( String brokerAddr, String topic, @@ -68,6 +73,15 @@ public class DefaultForwardClient extends AbstractForwardClient { return this.getClient().getMaxOffset(brokerAddr, topic, queueId, timeoutMillis); } + public CompletableFuture searchOffset( + String brokerAddr, + String topic, + int queueId, + long timestamp + ) { + return this.searchOffset(brokerAddr, topic, queueId, timestamp, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + } + public CompletableFuture searchOffset( String brokerAddr, String topic, 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 3ba34a6406..3752fd6f4c 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,6 +26,7 @@ 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.transaction.TransactionId; @@ -52,18 +53,21 @@ public class ForwardProducer extends AbstractForwardClient { return clientFactory.getTransactionalProducer(name, threadCount); } - public CompletableFuture heartBeat(String heartbeatAddr, HeartbeatData heartbeatData, long timeout) throws Exception { return this.getClient().sendHeartbeat(heartbeatAddr, heartbeatData, timeout); } public void endTransaction(String brokerAddr, EndTransactionRequestHeader requestHeader, long timeoutMillis) throws Exception { - this.getClient().endTransactionOneway( - brokerAddr, - requestHeader, - "end transaction from rmq proxy", - timeoutMillis - ); + this.getClient().endTransactionOneway(brokerAddr, requestHeader, "end transaction from rmq proxy", timeoutMillis); + } + + public CompletableFuture sendMessage( + String address, + String brokerName, + Message msg, + SendMessageRequestHeader requestHeader + ) { + return this.sendMessage(address, brokerName, msg, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture sendMessage( @@ -84,6 +88,10 @@ public class ForwardProducer extends AbstractForwardClient { }); } + public CompletableFuture sendMessageBack(String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader) { + return this.sendMessageBack(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + } + public CompletableFuture sendMessageBack(String brokerAddr, ConsumerSendMsgBackRequestHeader requestHeader, long timeoutMillis) { return this.getClient().sendMessageBack(brokerAddr, requestHeader, 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 869aeffc14..25e15c1397 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,8 +22,9 @@ 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.connector.factory.ForwardClientFactory; +import org.apache.rocketmq.proxy.common.utils.ProxyUtils; import org.apache.rocketmq.proxy.config.ConfigurationManager; +import org.apache.rocketmq.proxy.connector.factory.ForwardClientFactory; public class ForwardReadConsumer extends AbstractForwardClient { @@ -46,13 +47,21 @@ public class ForwardReadConsumer extends AbstractForwardClient { return clientFactory.getMQClient(name, threadCount); } - public CompletableFuture popMessage(String address, String brokerName, PopMessageRequestHeader requestHeader, - long timeoutMillis) { - return getClient().popMessage(address, brokerName, requestHeader, timeoutMillis); + public CompletableFuture popMessage( + String address, + String brokerName, + PopMessageRequestHeader requestHeader, + long timeoutMillis + ) { + return this.getClient().popMessage(address, brokerName, requestHeader, timeoutMillis); + } + + public CompletableFuture pullMessage(String address, PullMessageRequestHeader requestHeader) { + return this.pullMessage(address, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture pullMessage(String address, PullMessageRequestHeader requestHeader, long timeoutMillis) { - return getClient().pullMessage(address, requestHeader, timeoutMillis); + return this.getClient().pullMessage(address, requestHeader, timeoutMillis); } } 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 1c0d721cc0..082c8c5963 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,8 +22,9 @@ 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.connector.factory.ForwardClientFactory; +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.remoting.exception.RemotingException; public class ForwardWriteConsumer extends AbstractForwardClient { @@ -47,12 +48,24 @@ public class ForwardWriteConsumer extends AbstractForwardClient { return clientFactory.getMQClient(name, threadCount); } + public CompletableFuture ackMessage(String address, AckMessageRequestHeader requestHeader) { + return this.ackMessage(address, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + } + public CompletableFuture ackMessage( String address, AckMessageRequestHeader requestHeader, long timeoutMillis ) { - return getClient().ackMessage(address, requestHeader, timeoutMillis); + return this.getClient().ackMessage(address, requestHeader, timeoutMillis); + } + + public CompletableFuture changeInvisibleTimeAsync( + String address, + String brokerName, + ChangeInvisibleTimeRequestHeader requestHeader + ) { + return this.changeInvisibleTimeAsync(address, brokerName, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); } public CompletableFuture changeInvisibleTimeAsync( @@ -61,7 +74,7 @@ public class ForwardWriteConsumer extends AbstractForwardClient { ChangeInvisibleTimeRequestHeader requestHeader, long timeoutMillis ) { - return getClient().changeInvisibleTimeAsync(address, brokerName, requestHeader, timeoutMillis); + return this.getClient().changeInvisibleTimeAsync(address, brokerName, requestHeader, timeoutMillis); } public void updateConsumerOffsetOneWay( @@ -69,6 +82,6 @@ public class ForwardWriteConsumer extends AbstractForwardClient { UpdateConsumerOffsetRequestHeader header, long timeoutMillis ) throws RemotingException, InterruptedException { - getClient().updateConsumerOffsetOneWay(brokerAddr, header, timeoutMillis); + this.getClient().updateConsumerOffsetOneWay(brokerAddr, header, timeoutMillis); } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/route/MessageQueueSelector.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/route/MessageQueueSelector.java index 66be2134ff..eff85b1472 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/route/MessageQueueSelector.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/route/MessageQueueSelector.java @@ -153,10 +153,10 @@ public class MessageQueueSelector { } 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; + for (SelectableMessageQueue targetMessageQueue : queues) { + String queueBrokerName = targetMessageQueue.getBrokerName(); + if (queueBrokerName.equals(brokerName) && targetMessageQueue.getQueueId() == queueId) { + return targetMessageQueue; } } return null; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/route/SelectableMessageQueue.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/route/SelectableMessageQueue.java index f7b196629d..78a388214d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/route/SelectableMessageQueue.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/route/SelectableMessageQueue.java @@ -72,7 +72,7 @@ public class SelectableMessageQueue implements Comparable receiveMessage(Context ctx, ReceiveMessageRequest request) { - long timeRemaining = Context.current() - .getDeadline() - .timeRemaining(TimeUnit.MILLISECONDS); - long pollTime = timeRemaining - ConfigurationManager.getProxyConfig().getLongPollingReserveTimeInMillis(); - if (pollTime <= 0) { - pollTime = timeRemaining; - } + long pollTime = GrpcConverter.buildPollTimeFromContext(ctx); PopMessageRequestHeader requestHeader = GrpcConverter.buildPopMessageRequestHeader(request, pollTime); RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.POP_MESSAGE, requestHeader); command.makeCustomHeaderToNet(); @@ -392,13 +385,7 @@ public class LocalGrpcService extends AbstractStartAndShutdown implements GrpcFo @Override public CompletableFuture pullMessage(Context ctx, PullMessageRequest request) { - long timeRemaining = Context.current() - .getDeadline() - .timeRemaining(TimeUnit.MILLISECONDS); - long pollTime = timeRemaining - ConfigurationManager.getProxyConfig().getLongPollingReserveTimeInMillis(); - if (pollTime <= 0) { - pollTime = timeRemaining; - } + long pollTime = GrpcConverter.buildPollTimeFromContext(ctx); PullMessageRequestHeader requestHeader = GrpcConverter.buildPullMessageRequestHeader(request, pollTime); RemotingCommand command = RemotingCommand.createRequestCommand(RequestCode.PULL_MESSAGE, requestHeader); command.makeCustomHeaderToNet(); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/BaseService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/BaseService.java index 2c7bb94d3d..61857f4bf9 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/BaseService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/BaseService.java @@ -16,11 +16,14 @@ */ package org.apache.rocketmq.proxy.grpc.service.cluster; +import apache.rocketmq.v1.FilterExpression; +import apache.rocketmq.v1.Resource; import com.google.rpc.Code; import io.grpc.Context; import org.apache.commons.lang3.StringUtils; import org.apache.rocketmq.common.consumer.ReceiptHandle; import org.apache.rocketmq.proxy.connector.ConnectorManager; +import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.adapter.ProxyException; public class BaseService { @@ -49,4 +52,10 @@ public class BaseService { } return addr; } + + protected void checkSubscriptionData(Resource topic, FilterExpression filterExpression) { + // for checking filterExpression. + String topicName = GrpcConverter.wrapResourceWithNamespace(topic); + GrpcConverter.buildSubscriptionData(topicName, filterExpression); + } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java index 78c5d77736..c82d658b45 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ConsumerService.java @@ -25,6 +25,9 @@ import apache.rocketmq.v1.ReceiveMessageRequest; import apache.rocketmq.v1.ReceiveMessageResponse; import com.google.rpc.Code; import io.grpc.Context; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.client.consumer.AckResult; import org.apache.rocketmq.client.consumer.AckStatus; import org.apache.rocketmq.client.consumer.PopResult; @@ -37,38 +40,33 @@ import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHead import org.apache.rocketmq.common.protocol.header.PopMessageRequestHeader; import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; import org.apache.rocketmq.proxy.common.utils.FilterUtils; -import org.apache.rocketmq.proxy.common.utils.ProxyUtils; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.connector.ForwardProducer; import org.apache.rocketmq.proxy.connector.ForwardReadConsumer; import org.apache.rocketmq.proxy.connector.ForwardWriteConsumer; import org.apache.rocketmq.proxy.connector.route.SelectableMessageQueue; -import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.adapter.DelayPolicy; +import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.adapter.ProxyException; import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.adapter.ResponseHook; - -import java.util.ArrayList; -import java.util.List; -import java.util.concurrent.CompletableFuture; -import java.util.concurrent.TimeUnit; import org.apache.rocketmq.remoting.protocol.RemotingCommand; public class ConsumerService extends BaseService { - + private final DelayPolicy delayPolicy; private final ForwardReadConsumer readConsumer; private final ForwardWriteConsumer writeConsumer; + /** + * For sending messages back to broker. + */ private final ForwardProducer producer; private volatile ReadQueueSelector readQueueSelector; - private volatile ResponseHook receiveMessageHook = null; - private volatile ResponseHook ackNoMatchedMessageHook = null; - private volatile ResponseHook ackMessageHook = null; - private volatile ResponseHook nackMessageHook = null; - - private final DelayPolicy delayPolicy; + private volatile ResponseHook receiveMessageHook; + private volatile ResponseHook ackNoMatchedMessageHook; + private volatile ResponseHook ackMessageHook; + private volatile ResponseHook nackMessageHook; public ConsumerService(ConnectorManager connectorManager) { super(connectorManager); @@ -82,13 +80,15 @@ public class ConsumerService extends BaseService { public CompletableFuture receiveMessage(Context ctx, ReceiveMessageRequest request) { CompletableFuture future = new CompletableFuture<>(); + // register hook. future.whenComplete((response, throwable) -> { if (receiveMessageHook != null) { receiveMessageHook.beforeResponse(ctx, request, response, throwable); } }); + try { - PopMessageRequestHeader requestHeader = this.convertToPopMessageRequestHeader(ctx, request); + PopMessageRequestHeader requestHeader = this.buildPopMessageRequestHeader(ctx, request); SelectableMessageQueue messageQueue = this.readQueueSelector.select(ctx, request, requestHeader); if (messageQueue == null) { @@ -118,27 +118,19 @@ public class ConsumerService extends BaseService { return future; } - protected PopMessageRequestHeader convertToPopMessageRequestHeader(Context ctx, ReceiveMessageRequest request) { - // check filterExpression is correct or not - GrpcConverter.buildSubscriptionData(GrpcConverter.wrapResourceWithNamespace(request.getPartition().getTopic()), request.getFilterExpression()); - - long timeRemaining = ctx.getDeadline() - .timeRemaining(TimeUnit.MILLISECONDS); - long pollTime = timeRemaining - ConfigurationManager.getProxyConfig().getLongPollingReserveTimeInMillis(); - if (pollTime <= 0) { - pollTime = timeRemaining; - } - - return GrpcConverter.buildPopMessageRequestHeader(request, pollTime); + protected PopMessageRequestHeader buildPopMessageRequestHeader(Context ctx, ReceiveMessageRequest request) { + checkSubscriptionData(request.getPartition().getTopic(), request.getFilterExpression()); + return GrpcConverter.buildPopMessageRequestHeader(request, GrpcConverter.buildPollTimeFromContext(ctx)); } protected ReceiveMessageResponse convertToReceiveMessageResponse(Context ctx, ReceiveMessageRequest request, PopResult result) { - SubscriptionData subscriptionData = GrpcConverter.buildSubscriptionData( - GrpcConverter.wrapResourceWithNamespace(request.getPartition().getTopic()), request.getFilterExpression()); PopStatus status = result.getPopStatus(); switch (status) { case FOUND: - break; + return ReceiveMessageResponse.newBuilder() + .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) + .addAllMessages(checkAndGetMessagesFromPopResult(ctx, request, result)) + .build(); case POLLING_FULL: return ReceiveMessageResponse.newBuilder() .setCommon(ResponseBuilder.buildCommon(Code.RESOURCE_EXHAUSTED, "polling full")) @@ -150,6 +142,11 @@ public class ConsumerService extends BaseService { .setCommon(ResponseBuilder.buildCommon(Code.OK, "no new message")) .build(); } + } + + protected List checkAndGetMessagesFromPopResult(Context ctx, ReceiveMessageRequest request, PopResult result) { + String topicName = GrpcConverter.wrapResourceWithNamespace(request.getPartition().getTopic()); + SubscriptionData subscriptionData = GrpcConverter.buildSubscriptionData(topicName, request.getFilterExpression()); List messages = new ArrayList<>(); for (MessageExt messageExt : result.getMsgFoundList()) { @@ -160,14 +157,12 @@ public class ConsumerService extends BaseService { messages.add(GrpcConverter.buildMessage(messageExt)); } - return ReceiveMessageResponse.newBuilder() - .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) - .addAllMessages(messages) - .build(); + return messages; } protected void ackNoMatchedMessage(Context ctx, ReceiveMessageRequest request, MessageExt messageExt) { CompletableFuture future = new CompletableFuture<>(); + AckMessageRequestHeader ackMessageRequestHeader = new AckMessageRequestHeader(); try { ReceiptHandle handle = ReceiptHandle.create(messageExt); @@ -181,15 +176,17 @@ public class ConsumerService extends BaseService { ackMessageRequestHeader.setExtraInfo(handle.getReceiptHandle()); ackMessageRequestHeader.setOffset(handle.getOffset()); - future = this.writeConsumer.ackMessage(brokerAddr, ackMessageRequestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + future = this.writeConsumer.ackMessage(brokerAddr, ackMessageRequestHeader); } catch (Throwable t) { future.completeExceptionally(t); } + future.whenComplete((ackResult, throwable) -> { if (ackNoMatchedMessageHook != null) { ackNoMatchedMessageHook.beforeResponse(ctx, ackMessageRequestHeader, ackResult, throwable); } }); + } public CompletableFuture ackMessage(Context ctx, AckMessageRequest request) { @@ -199,31 +196,32 @@ public class ConsumerService extends BaseService { ackMessageHook.beforeResponse(ctx, request, response, throwable); } }); + try { ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle()); String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName()); - AckMessageRequestHeader requestHeader = this.convertToAckMessageRequestHeader(ctx, request); - CompletableFuture ackResultFuture = this.writeConsumer.ackMessage(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + AckMessageRequestHeader requestHeader = this.buildAckMessageRequestHeader(ctx, request); + CompletableFuture ackResultFuture = this.writeConsumer.ackMessage(brokerAddr, requestHeader); ackResultFuture - .thenAccept(result -> { - try { - future.complete(convertToAckMessageResponse(ctx, request, result)); - } catch (Throwable throwable) { - future.completeExceptionally(throwable); - } - }) - .exceptionally(throwable -> { + .thenAccept(result -> { + try { + future.complete(convertToAckMessageResponse(ctx, request, result)); + } catch (Throwable throwable) { future.completeExceptionally(throwable); - return null; - }); + } + }) + .exceptionally(throwable -> { + future.completeExceptionally(throwable); + return null; + }); } catch (Throwable t) { future.completeExceptionally(t); } return future; } - protected AckMessageRequestHeader convertToAckMessageRequestHeader(Context ctx, AckMessageRequest request) { + protected AckMessageRequestHeader buildAckMessageRequestHeader(Context ctx, AckMessageRequest request) { return GrpcConverter.buildAckMessageRequestHeader(request); } @@ -252,8 +250,9 @@ public class ConsumerService extends BaseService { if (request.getDeliveryAttempt() >= request.getMaxDeliveryAttempts()) { CompletableFuture resultFuture = this.producer.sendMessageBack( brokerAddr, - this.convertToConsumerSendMsgBackToDLQRequestHeader(ctx, request), - ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + this.buildConsumerSendMsgBackToDLQRequestHeader(ctx, request) + ); + resultFuture .thenAccept(result -> { try { @@ -267,9 +266,8 @@ public class ConsumerService extends BaseService { return null; }); } else { - ChangeInvisibleTimeRequestHeader requestHeader = this.convertToChangeInvisibleTimeRequestHeader(ctx, request); - CompletableFuture resultFuture = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader, - ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + ChangeInvisibleTimeRequestHeader requestHeader = this.buildChangeInvisibleTimeRequestHeader(ctx, request); + CompletableFuture resultFuture = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader); resultFuture .thenAccept(result -> { try { @@ -289,11 +287,11 @@ public class ConsumerService extends BaseService { return future; } - protected ChangeInvisibleTimeRequestHeader convertToChangeInvisibleTimeRequestHeader(Context ctx, NackMessageRequest request) { - return GrpcConverter.buildChangeInvisibleTimeRequestHeader(request, delayPolicy); + protected ChangeInvisibleTimeRequestHeader buildChangeInvisibleTimeRequestHeader(Context ctx, NackMessageRequest request) { + return GrpcConverter.buildChangeInvisibleTimeRequestHeader(request, this.delayPolicy); } - protected ConsumerSendMsgBackRequestHeader convertToConsumerSendMsgBackToDLQRequestHeader(Context ctx, NackMessageRequest request) { + protected ConsumerSendMsgBackRequestHeader buildConsumerSendMsgBackToDLQRequestHeader(Context ctx, NackMessageRequest request) { return GrpcConverter.buildConsumerSendMsgBackToDLQRequestHeader(request); } @@ -318,23 +316,19 @@ public class ConsumerService extends BaseService { this.readQueueSelector = readQueueSelector; } - public void setReceiveMessageHook( - ResponseHook receiveMessageHook) { + public void setReceiveMessageHook(ResponseHook receiveMessageHook) { this.receiveMessageHook = receiveMessageHook; } - public void setAckNoMatchedMessageHook( - ResponseHook ackNoMatchedMessageHook) { + public void setAckNoMatchedMessageHook(ResponseHook ackNoMatchedMessageHook) { this.ackNoMatchedMessageHook = ackNoMatchedMessageHook; } - public void setAckMessageHook( - ResponseHook ackMessageHook) { + public void setAckMessageHook(ResponseHook ackMessageHook) { this.ackMessageHook = ackMessageHook; } - public void setNackMessageHook( - ResponseHook nackMessageHook) { + public void setNackMessageHook(ResponseHook nackMessageHook) { this.nackMessageHook = nackMessageHook; } } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultWriteQueueSelector.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultWriteQueueSelector.java index 1cc9ba7bc6..19948a0c5a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultWriteQueueSelector.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/DefaultWriteQueueSelector.java @@ -27,8 +27,8 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; public class DefaultWriteQueueSelector implements WriteQueueSelector { + private static final Logger LOGGER = LoggerFactory.getLogger(DefaultWriteQueueSelector.class); - private static final Logger log = LoggerFactory.getLogger(DefaultWriteQueueSelector.class); protected final TopicRouteCache topicRouteCache; public DefaultWriteQueueSelector(TopicRouteCache topicRouteCache) { @@ -36,9 +36,12 @@ public class DefaultWriteQueueSelector implements WriteQueueSelector { } @Override - public SelectableMessageQueue selectQueue(Context ctx, SendMessageRequest request, + public SelectableMessageQueue selectQueue( + Context ctx, + SendMessageRequest request, SendMessageRequestHeader requestHeader, - org.apache.rocketmq.common.message.Message message) { + org.apache.rocketmq.common.message.Message message + ) { try { String topic = requestHeader.getTopic(); String brokerName = ""; @@ -47,19 +50,19 @@ public class DefaultWriteQueueSelector implements WriteQueueSelector { } Integer queueId = requestHeader.getQueueId(); String shardingKey = message.getProperty(MessageConst.PROPERTY_SHARDING_KEY); - SelectableMessageQueue addressableMessageQueue; - if (!StringUtils.isBlank(brokerName) && queueId != null) { + SelectableMessageQueue targetMessageQueue; + if (StringUtils.isNotBlank(brokerName) && queueId != null) { // Grpc client sendSelect situation - addressableMessageQueue = selectTargetQueue(topic, brokerName, queueId); + targetMessageQueue = selectTargetQueue(topic, brokerName, queueId); } else if (shardingKey != null) { // With shardingKey - addressableMessageQueue = selectOrderQueue(topic, shardingKey); + targetMessageQueue = selectOrderQueue(topic, shardingKey); } else { - addressableMessageQueue = selectNormalQueue(topic); + targetMessageQueue = selectNormalQueue(topic); } - return addressableMessageQueue; + return targetMessageQueue; } catch (Exception e) { - log.error("error when select queue in DefaultMessageQueueSelector. request: {}", request, e); + LOGGER.error("error when select queue in DefaultMessageQueueSelector. request: {}", request, e); return null; } } @@ -68,8 +71,7 @@ public class DefaultWriteQueueSelector implements WriteQueueSelector { return this.topicRouteCache.selectOneWriteQueue(topic, null); } - protected SelectableMessageQueue selectTargetQueue(String topic, String brokerName, - int queueId) throws Exception { + protected SelectableMessageQueue selectTargetQueue(String topic, String brokerName, int queueId) throws Exception { return this.topicRouteCache.selectOneWriteQueue(topic, brokerName, queueId); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerService.java index e924e28855..1e23ce167d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ProducerService.java @@ -31,8 +31,8 @@ import org.apache.rocketmq.client.producer.SendStatus; import org.apache.rocketmq.common.consumer.ReceiptHandle; import org.apache.rocketmq.common.protocol.header.ConsumerSendMsgBackRequestHeader; import org.apache.rocketmq.common.protocol.header.SendMessageRequestHeader; -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.route.SelectableMessageQueue; import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.adapter.ProxyException; @@ -42,12 +42,15 @@ import org.apache.rocketmq.remoting.protocol.RemotingCommand; public class ProducerService extends BaseService { + private final ForwardProducer producer; private volatile WriteQueueSelector writeQueueSelector; - private volatile ResponseHook sendMessageHook = null; - private volatile ResponseHook forwardMessageToDLQHook = null; + private volatile ResponseHook sendMessageHook; + private volatile ResponseHook forwardMessageToDLQHook; + public ProducerService(ConnectorManager connectorManager) { super(connectorManager); + this.producer = connectorManager.getForwardProducer(); writeQueueSelector = new DefaultWriteQueueSelector(this.connectorManager.getTopicRouteCache()); } @@ -73,23 +76,24 @@ public class ProducerService extends BaseService { }); try { - Pair requestPair = this.convertSendMessageRequest(ctx, request); + Pair requestPair = this.buildSendMessageRequest(ctx, request); SendMessageRequestHeader requestHeader = requestPair.getLeft(); org.apache.rocketmq.common.message.Message message = requestPair.getRight(); - SelectableMessageQueue addressableMessageQueue = writeQueueSelector.selectQueue(ctx, request, requestHeader, message); + SelectableMessageQueue selectableMessageQueue = writeQueueSelector.selectQueue(ctx, request, requestHeader, message); String topic = requestHeader.getTopic(); - if (addressableMessageQueue == null) { - throw new ProxyException(Code.NOT_FOUND, "no writeable topic route for topic " + topic); + if (selectableMessageQueue == null) { + throw new ProxyException(Code.NOT_FOUND, "no writeable topic route for topic: " + topic); } - CompletableFuture sendResultCompletableFuture = this.connectorManager.getForwardProducer().sendMessage( - addressableMessageQueue.getBrokerAddr(), - addressableMessageQueue.getBrokerName(), + // send message to broker. + CompletableFuture sendResultCompletableFuture = this.producer.sendMessage( + selectableMessageQueue.getBrokerAddr(), + selectableMessageQueue.getBrokerName(), message, - requestHeader, - ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT + requestHeader ); + sendResultCompletableFuture .thenAccept(result -> { try { @@ -108,20 +112,21 @@ public class ProducerService extends BaseService { return future; } - protected Pair convertSendMessageRequest( + protected Pair buildSendMessageRequest( Context ctx, SendMessageRequest request) { - return Pair.of(GrpcConverter.buildSendMessageRequestHeader(request), GrpcConverter.buildMessage(request.getMessage())); + SendMessageRequestHeader requestHeader = GrpcConverter.buildSendMessageRequestHeader(request); + org.apache.rocketmq.common.message.Message message = GrpcConverter.buildMessage(request.getMessage()); + return Pair.of(requestHeader, message); } - protected SendMessageResponse convertToSendMessageResponse(Context ctx, SendMessageRequest request, - SendResult sendResult) { - if (sendResult.getSendStatus() != SendStatus.SEND_OK) { + protected SendMessageResponse convertToSendMessageResponse(Context ctx, SendMessageRequest request, SendResult result) { + if (result.getSendStatus() != SendStatus.SEND_OK) { return SendMessageResponse.newBuilder() - .setCommon(ResponseBuilder.buildCommon(Code.INTERNAL, "send message failed, sendStatus=" + sendResult.getSendStatus())) + .setCommon(ResponseBuilder.buildCommon(Code.INTERNAL, "send message failed, sendStatus=" + result.getSendStatus())) .build(); } - if (StringUtils.isNotBlank(sendResult.getTransactionId())) { + if (StringUtils.isNotBlank(result.getTransactionId())) { Message message = request.getMessage(); String group = GrpcConverter.wrapResourceWithNamespace(message.getSystemAttribute().getProducerGroup()); String topic = GrpcConverter.wrapResourceWithNamespace(message.getTopic()); @@ -130,8 +135,8 @@ public class ProducerService extends BaseService { return SendMessageResponse.newBuilder() .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) - .setMessageId(StringUtils.defaultString(sendResult.getMsgId())) - .setTransactionId(StringUtils.defaultString(sendResult.getTransactionId())) + .setMessageId(StringUtils.defaultString(result.getMsgId())) + .setTransactionId(StringUtils.defaultString(result.getTransactionId())) // use "" if transactionID is null. .build(); } @@ -143,12 +148,12 @@ public class ProducerService extends BaseService { forwardMessageToDLQHook.beforeResponse(ctx, request, response, throwable); } }); + try { ReceiptHandle receiptHandle = this.resolveReceiptHandle(ctx, request.getReceiptHandle()); String brokerAddr = this.getBrokerAddr(ctx, receiptHandle.getBrokerName()); - ConsumerSendMsgBackRequestHeader requestHeader = this.convertToConsumerSendMsgBackRequestHeader(ctx, request); - CompletableFuture resultFuture = this.connectorManager.getForwardProducer() - .sendMessageBack(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + ConsumerSendMsgBackRequestHeader requestHeader = this.buildConsumerSendMsgBackRequestHeader(ctx, request); + CompletableFuture resultFuture = this.producer.sendMessageBack(brokerAddr, requestHeader); resultFuture .thenAccept(result -> future.complete( @@ -167,7 +172,7 @@ public class ProducerService extends BaseService { return future; } - protected ConsumerSendMsgBackRequestHeader convertToConsumerSendMsgBackRequestHeader(Context ctx, + protected ConsumerSendMsgBackRequestHeader buildConsumerSendMsgBackRequestHeader(Context ctx, ForwardMessageToDeadLetterQueueRequest request) { return GrpcConverter.buildConsumerSendMsgBackRequestHeader(request); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/PullMessageService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/PullMessageService.java index 4193609f94..b3fc812089 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/PullMessageService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/PullMessageService.java @@ -20,7 +20,6 @@ import apache.rocketmq.v1.Message; import apache.rocketmq.v1.Partition; import apache.rocketmq.v1.PullMessageRequest; import apache.rocketmq.v1.PullMessageResponse; -import apache.rocketmq.v1.QueryOffsetPolicy; import apache.rocketmq.v1.QueryOffsetRequest; import apache.rocketmq.v1.QueryOffsetResponse; import com.google.protobuf.util.Timestamps; @@ -28,32 +27,31 @@ import com.google.rpc.Code; import io.grpc.Context; import java.util.List; import java.util.concurrent.CompletableFuture; -import java.util.concurrent.TimeUnit; import java.util.stream.Collectors; import org.apache.rocketmq.client.consumer.PullResult; import org.apache.rocketmq.client.consumer.PullStatus; import org.apache.rocketmq.common.protocol.header.PullMessageRequestHeader; import org.apache.rocketmq.common.protocol.heartbeat.SubscriptionData; import org.apache.rocketmq.proxy.common.utils.FilterUtils; -import org.apache.rocketmq.proxy.common.utils.ProxyUtils; -import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.connector.ConnectorManager; import org.apache.rocketmq.proxy.connector.DefaultForwardClient; +import org.apache.rocketmq.proxy.connector.ForwardReadConsumer; import org.apache.rocketmq.proxy.grpc.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.adapter.ResponseHook; public class PullMessageService extends BaseService { - private final DefaultForwardClient defaultForwardClient; + private final DefaultForwardClient forwardClient; + private final ForwardReadConsumer readConsumer; - private volatile ResponseHook queryOffsetHook = null; - - private volatile ResponseHook pullMessageHook = null; + private volatile ResponseHook queryOffsetHook; + private volatile ResponseHook pullMessageHook; public PullMessageService(ConnectorManager connectorManager) { super(connectorManager); - this.defaultForwardClient = connectorManager.getDefaultForwardClient(); + this.forwardClient = connectorManager.getDefaultForwardClient(); + this.readConsumer = connectorManager.getForwardReadConsumer(); } public CompletableFuture queryOffset(Context ctx, QueryOffsetRequest request) { @@ -63,24 +61,28 @@ public class PullMessageService extends BaseService { queryOffsetHook.beforeResponse(ctx, request, response, throwable); } }); + try { Partition partition = request.getPartition(); String topic = GrpcConverter.wrapResourceWithNamespace(partition.getTopic()); - String brokerName = partition.getBroker().getName(); int queueId = partition.getId(); + CompletableFuture offsetFuture; - if (request.getPolicy() == QueryOffsetPolicy.BEGINNING) { - offsetFuture = CompletableFuture.completedFuture(0L); - } else if (request.getPolicy() == QueryOffsetPolicy.END) { - String brokerAddr = this.getBrokerAddr(ctx, brokerName); - offsetFuture = this.defaultForwardClient.getMaxOffset(brokerAddr, topic, queueId, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); - } else { - long timestamp = Timestamps.toMillis(request.getTimePoint()); - String brokerAddr = this.getBrokerAddr(ctx, brokerName); - offsetFuture = this.defaultForwardClient.searchOffset(brokerAddr, topic, queueId, timestamp, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + switch (request.getPolicy()) { + case BEGINNING: + offsetFuture = CompletableFuture.completedFuture(0L); + break; + case END: + offsetFuture = this.forwardClient.getMaxOffset(this.getBrokerAddr(ctx, brokerName), topic, queueId); + break; + default: + long timestamp = Timestamps.toMillis(request.getTimePoint()); + offsetFuture = this.forwardClient.searchOffset(this.getBrokerAddr(ctx, brokerName), topic, queueId, timestamp); } - offsetFuture.thenAccept(result -> future.complete( + + offsetFuture + .thenAccept(result -> future.complete( QueryOffsetResponse.newBuilder() .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) .setOffset(result) @@ -104,13 +106,12 @@ public class PullMessageService extends BaseService { }); try { - PullMessageRequestHeader requestHeader = this.convertToPullMessageRequestHeader(ctx, request); + PullMessageRequestHeader requestHeader = this.buildPullMessageRequestHeader(ctx, request); String brokerName = request.getPartition().getBroker().getName(); String brokerAddr = this.getBrokerAddr(ctx, brokerName); - CompletableFuture pullResultFuture = this.connectorManager.getForwardReadConsumer() - .pullMessage(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + CompletableFuture pullResultFuture = this.readConsumer.pullMessage(brokerAddr, requestHeader); pullResultFuture .thenAccept(pullResult -> { try { @@ -130,17 +131,9 @@ public class PullMessageService extends BaseService { return future; } - protected PullMessageRequestHeader convertToPullMessageRequestHeader(Context ctx, PullMessageRequest request) { - // check filterExpression is correct or not - GrpcConverter.buildSubscriptionData(GrpcConverter.wrapResourceWithNamespace(request.getPartition().getTopic()), request.getFilterExpression()); - - long timeRemaining = ctx.getDeadline() - .timeRemaining(TimeUnit.MILLISECONDS); - long pollTime = timeRemaining - ConfigurationManager.getProxyConfig().getLongPollingReserveTimeInMillis(); - if (pollTime <= 0) { - pollTime = timeRemaining; - } - return GrpcConverter.buildPullMessageRequestHeader(request, pollTime); + protected PullMessageRequestHeader buildPullMessageRequestHeader(Context ctx, PullMessageRequest request) { + checkSubscriptionData(request.getPartition().getTopic(), request.getFilterExpression()); + return GrpcConverter.buildPullMessageRequestHeader(request, GrpcConverter.buildPollTimeFromContext(ctx)); } protected PullMessageResponse convertToPullMessageResponse(Context ctx, PullMessageRequest request, PullResult result) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java index 7ec10d6b5a..5b48a15d4a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/RouteService.java @@ -56,11 +56,11 @@ public class RouteService extends BaseService { private final ProxyMode mode; private volatile ParameterConverter queryRouteEndpointConverter; - private volatile ResponseHook queryRouteHook = null; + private volatile ResponseHook queryRouteHook; private volatile ParameterConverter queryAssignmentEndpointConverter; private volatile AssignmentQueueSelector assignmentQueueSelector; - private volatile ResponseHook queryAssignmentHook = null; + private volatile ResponseHook queryAssignmentHook; public RouteService(ProxyMode mode, ConnectorManager connectorManager) { super(connectorManager); @@ -79,8 +79,7 @@ public class RouteService extends BaseService { this.queryRouteHook = queryRouteHook; } - public void setQueryAssignmentEndpointConverter( - ParameterConverter queryAssignmentEndpointConverter) { + public void setQueryAssignmentEndpointConverter(ParameterConverter queryAssignmentEndpointConverter) { this.queryAssignmentEndpointConverter = queryAssignmentEndpointConverter; } @@ -88,8 +87,7 @@ public class RouteService extends BaseService { this.assignmentQueueSelector = assignmentQueueSelector; } - public void setQueryAssignmentHook( - ResponseHook queryAssignmentHook) { + public void setQueryAssignmentHook(ResponseHook queryAssignmentHook) { this.queryAssignmentHook = queryAssignmentHook; } @@ -102,8 +100,8 @@ public class RouteService extends BaseService { }); try { - MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache() - .getMessageQueue(GrpcConverter.wrapResourceWithNamespace(request.getTopic())); + String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); + MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache().getMessageQueue(topicName); TopicRouteData topicRouteData = messageQueueWrapper.getTopicRouteData(); List queueDataList = topicRouteData.getQueueDatas(); List brokerDataList = topicRouteData.getBrokerDatas(); @@ -217,8 +215,8 @@ public class RouteService extends BaseService { List assignments = new ArrayList<>(); List messageQueueList = this.assignmentQueueSelector.getAssignment(ctx, request); if (ProxyMode.isLocalMode(mode)) { - MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache() - .getMessageQueue(GrpcConverter.wrapResourceWithNamespace(request.getTopic())); + String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); + MessageQueueWrapper messageQueueWrapper = this.connectorManager.getTopicRouteCache().getMessageQueue(topicName); TopicRouteData topicRouteData = messageQueueWrapper.getTopicRouteData(); Map> brokerMap = buildBrokerMap(topicRouteData.getBrokerDatas()); for (SelectableMessageQueue messageQueue : messageQueueList) { 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 f0f730f17e..bcc4796d8f 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 @@ -46,8 +46,8 @@ public class TransactionService extends BaseService implements TransactionStateC private final ChannelManager channelManager; private final ForwardProducer forwardProducer; - private volatile ResponseHook checkTransactionStateHook = null; - private volatile ResponseHook endTransactionHook = null; + private volatile ResponseHook checkTransactionStateHook; + private volatile ResponseHook endTransactionHook; public TransactionService(ConnectorManager connectorManager, ChannelManager channelManager) { super(connectorManager); @@ -63,11 +63,11 @@ public class TransactionService extends BaseService implements TransactionStateC if (CollectionUtils.isEmpty(clientIdList)) { return; } + String clientId = clientIdList.get(ThreadLocalRandom.current().nextInt(clientIdList.size())); - GrpcClientChannel channel = GrpcClientChannel.getChannel(this.channelManager, checkData.getGroupId(), clientId); - String transactionId = checkData.getTransactionId().getProxyTransactionId(); + String transactionId = checkData.getTransactionId().getProxyTransactionId(); Message message = GrpcConverter.buildMessage(checkData.getMessageExt()); PollCommandResponse response = PollCommandResponse.newBuilder() .setRecoverOrphanedTransactionCommand( @@ -76,6 +76,7 @@ public class TransactionService extends BaseService implements TransactionStateC .setTransactionId(transactionId) .build() ).build(); + channel.writeAndFlush(response); if (this.checkTransactionStateHook != null) { this.checkTransactionStateHook.beforeResponse(ctx, checkData, response, null); @@ -97,10 +98,9 @@ public class TransactionService extends BaseService implements TransactionStateC try { TransactionId handle = TransactionId.decode(request.getTransactionId()); + String brokerAddr = RemotingHelper.parseSocketAddressAddr(handle.getBrokerAddr()); EndTransactionRequestHeader requestHeader = this.toEndTransactionRequestHeader(ctx, request); - this.forwardProducer.endTransaction( - RemotingHelper.parseSocketAddressAddr(handle.getBrokerAddr()), - requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + this.forwardProducer.endTransaction(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); future.complete(EndTransactionResponse.newBuilder() .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) .build());