From c08fee78affef8a3ccd09bc8900aef1f2721a5ba Mon Sep 17 00:00:00 2001 From: "Jixiang.jjx" Date: Sun, 24 Apr 2022 19:11:03 +0800 Subject: [PATCH] [ISSUE #3949] do the code refactoring work for readability. --- .../java/org/apache/rocketmq/proxy/ProxyStartup.java | 2 +- .../rocketmq/proxy/grpc/v2/adapter/GrpcConverter.java | 11 ++++++----- .../grpc/v2/service/cluster/ProducerService.java | 2 +- 3 files changed, 8 insertions(+), 7 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 55f7526e58..516a99506a 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/ProxyStartup.java @@ -81,7 +81,7 @@ public class ProxyStartup { System.exit(1); } - System.out.printf("%s%n", new Date() + " rmq-proxy startup successfully"); + System.out.println(new Date() + " rmq-proxy startup successfully"); log.info(new Date() + " rmq-proxy startup successfully"); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/GrpcConverter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/GrpcConverter.java index 9425d5b64f..d933b5e27d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/GrpcConverter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/adapter/GrpcConverter.java @@ -89,7 +89,6 @@ import org.apache.rocketmq.common.sysflag.MessageSysFlag; import org.apache.rocketmq.common.utils.BinaryUtil; import org.apache.rocketmq.logging.InternalLogger; import org.apache.rocketmq.logging.InternalLoggerFactory; -import org.apache.rocketmq.proxy.common.DelayPolicy; import org.apache.rocketmq.proxy.common.utils.ProxyUtils; import org.apache.rocketmq.proxy.config.ConfigurationManager; import org.apache.rocketmq.proxy.connector.transaction.TransactionId; @@ -444,22 +443,24 @@ public class GrpcConverter { } public static List buildMessage(List protoMessageList, - Resource topic, String producerGroup) { + Resource topic) { + String topicName = wrapResourceWithNamespace(topic); List messages = new ArrayList<>(); for (Message protoMessage : protoMessageList) { if (!protoMessage.getTopic().equals(topic)) { throw new ProxyException(Code.MESSAGE_CORRUPTED, "topic in message is not same"); } - messages.add(buildMessage(protoMessage, producerGroup)); + // here use topicName as producerGroup for transactional checker. + messages.add(buildMessage(protoMessage, topicName)); } return messages; } public static org.apache.rocketmq.common.message.Message buildMessage(Message protoMessage, String producerGroup) { - String topic = wrapResourceWithNamespace(protoMessage.getTopic()); + String topicName = wrapResourceWithNamespace(protoMessage.getTopic()); org.apache.rocketmq.common.message.Message message = - new org.apache.rocketmq.common.message.Message(topic, protoMessage.getBody().toByteArray()); + new org.apache.rocketmq.common.message.Message(topicName, protoMessage.getBody().toByteArray()); Map messageProperty = buildMessageProperty(protoMessage, producerGroup); MessageAccessor.setProperties(message, messageProperty); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerService.java index 3bd711c194..39136e36d6 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ProducerService.java @@ -108,7 +108,7 @@ public class ProducerService extends BaseService { // use topic name as group Resource topic = request.getMessages(0).getTopic(); String topicName = GrpcConverter.wrapResourceWithNamespace(topic); - return GrpcConverter.buildMessage(request.getMessagesList(), topic, topicName); + return GrpcConverter.buildMessage(request.getMessagesList(), topic); } protected SendMessageResponse convertToSendMessageResponse(Context ctx, SendMessageRequest request,