From ff3117f322418bf829718ddf657d0d5865a94a66 Mon Sep 17 00:00:00 2001 From: "Jixiang.jjx" Date: Mon, 28 Mar 2022 14:15:49 +0800 Subject: [PATCH] [ISSUE #3949] Do some refactoring work. --- .../proxy/connector/AbstractForwardClient.java | 3 ++- .../rocketmq/proxy/connector/ForwardProducer.java | 2 +- .../connector/factory/ForwardClientManager.java | 2 +- .../proxy/connector/factory/MQClientFactory.java | 3 +-- .../processor/ProxyClientRemotingProcessor.java | 2 +- .../proxy/connector/route/TopicRouteCache.java | 3 ++- .../TransactionHeartbeatRegisterService.java | 1 - .../connector/transaction/TransactionId.java | 6 +++--- .../proxy/grpc/GrpcMessagingProcessor.java | 15 ++++++++++----- .../apache/rocketmq/proxy/grpc/GrpcServer.java | 5 +++-- .../proxy/grpc/adapter/GrpcConverter.java | 15 +++++---------- .../service/cluster/ForwardClientService.java | 6 ++++-- .../connector/transaction/TransactionIdTest.java | 6 +++--- .../proxy/grpc/service/LocalGrpcServiceTest.java | 2 +- .../service/cluster/TransactionServiceTest.java | 4 ++-- 15 files changed, 39 insertions(+), 36 deletions(-) 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 c8ca3f563c..85ad25d855 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 @@ -46,7 +46,8 @@ public abstract class AbstractForwardClient implements StartAndShutdown { if (clients.length == 1) { return this.clients[0]; } - return this.clients[ThreadLocalRandom.current().nextInt(this.clients.length)]; + int index = ThreadLocalRandom.current().nextInt(this.clients.length); + return this.clients[index]; } @Override 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 5c362713e9..f31fe95c32 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 @@ -87,7 +87,7 @@ public class ForwardProducer extends AbstractForwardClient { return future.thenApply(sendResult -> { int tranType = MessageSysFlag.getTransactionValue(requestHeader.getSysFlag()); if (SendStatus.SEND_OK.equals(sendResult.getSendStatus()) && tranType == MessageSysFlag.TRANSACTION_PREPARED_TYPE) { - TransactionId transactionId = TransactionId.genFromBrokerTransactionId(address, sendResult); + TransactionId transactionId = TransactionId.genByBrokerTransactionId(address, sendResult); sendResult.setTransactionId(transactionId.getProxyTransactionId()); } return sendResult; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/ForwardClientManager.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/ForwardClientManager.java index 37ac2f2497..0aaa079581 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/ForwardClientManager.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/ForwardClientManager.java @@ -30,7 +30,7 @@ import org.apache.rocketmq.remoting.RPCHook; public class ForwardClientManager implements StartAndShutdown { - private RPCHook rpcHook = null; + private RPCHook rpcHook; private final MQClientFactory mqClientFactory; private final TransactionProducerFactory transactionalProducerFactory; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/MQClientFactory.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/MQClientFactory.java index 07484bd54b..450c611bd7 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/MQClientFactory.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/factory/MQClientFactory.java @@ -23,8 +23,7 @@ import org.apache.rocketmq.remoting.RPCHook; public class MQClientFactory extends AbstractMQClientFactory { - public MQClientFactory(ScheduledExecutorService scheduledExecutorService, - RPCHook rpcHook) { + public MQClientFactory(ScheduledExecutorService scheduledExecutorService, RPCHook rpcHook) { super(scheduledExecutorService, rpcHook); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/processor/ProxyClientRemotingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/processor/ProxyClientRemotingProcessor.java index 3f217294f4..41001e8f6e 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/processor/ProxyClientRemotingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/processor/ProxyClientRemotingProcessor.java @@ -63,7 +63,7 @@ public class ProxyClientRemotingProcessor extends ClientRemotingProcessor { requestHeader.getTranStateTableOffset(), requestHeader.getCommitLogOffset(), requestHeader.getMsgId(), - TransactionId.genFromBrokerTransactionId( + TransactionId.genByBrokerTransactionId( ctx.channel().remoteAddress(), requestHeader.getTransactionId(), requestHeader.getCommitLogOffset(), 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 49189a0126..ad4ed98106 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 @@ -81,7 +81,8 @@ public class TopicRouteCache { } public SelectableMessageQueue selectOneWriteQueue(String topic, String brokerName, int queueId) throws Exception { - return getMessageQueue(topic).getWriteSelector().selectOne(brokerName, queueId); + return getMessageQueue(topic).getWriteSelector() + .selectOne(brokerName, queueId); } public SelectableMessageQueue selectOneWriteQueueByKey(String topic, String shardingKey) throws Exception { 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 d7f543b49e..a54cb31a39 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 @@ -40,7 +40,6 @@ import org.slf4j.Logger; import org.slf4j.LoggerFactory; public class TransactionHeartbeatRegisterService implements StartAndShutdown { - private static final Logger log = LoggerFactory.getLogger(TransactionHeartbeatRegisterService.class); private static final String TRANS_HEARTBEAT_CLIENT_ID = "rmq-proxy-producer-client"; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionId.java b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionId.java index b909fbecd3..387bb5800d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionId.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/connector/transaction/TransactionId.java @@ -54,7 +54,7 @@ public class TransactionId { public TransactionId() { } - public static TransactionId genFromBrokerTransactionId(String brokerAddr, SendResult sendResult) { + public static TransactionId genByBrokerTransactionId(String brokerAddr, SendResult sendResult) { MessageId id = new MessageId(null, 0); try { if (sendResult.getOffsetMsgId() != null) { @@ -65,11 +65,11 @@ public class TransactionId { } catch (Exception e) { log.warn("genFromBrokerTransactionId failed. brokerAddr: {}, sendResult: {}", brokerAddr, sendResult, e); } - return genFromBrokerTransactionId(RemotingUtil.string2SocketAddress(brokerAddr), sendResult.getTransactionId(), + return genByBrokerTransactionId(RemotingUtil.string2SocketAddress(brokerAddr), sendResult.getTransactionId(), id.getOffset(), sendResult.getQueueOffset()); } - public static TransactionId genFromBrokerTransactionId(SocketAddress brokerAddr, String orgTransactionId, + public static TransactionId genByBrokerTransactionId(SocketAddress brokerAddr, String orgTransactionId, long commitLogOffset, long tranStateTableOffset) { byte[] orgTransactionIdByte = new byte[0]; if (StringUtils.isNotBlank(orgTransactionId)) { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcMessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcMessagingProcessor.java index f520cc36be..fe8a4f6aac 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcMessagingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcMessagingProcessor.java @@ -189,7 +189,8 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic } @Override - public void forwardMessageToDeadLetterQueue(ForwardMessageToDeadLetterQueueRequest request, StreamObserver responseObserver) { + public void forwardMessageToDeadLetterQueue(ForwardMessageToDeadLetterQueueRequest request, + StreamObserver responseObserver) { CompletableFuture future = grpcForwardService.forwardMessageToDeadLetterQueue(Context.current(), request); future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) .exceptionally(e -> { @@ -251,7 +252,8 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic } @Override - public void reportThreadStackTrace(ReportThreadStackTraceRequest request, StreamObserver responseObserver) { + public void reportThreadStackTrace(ReportThreadStackTraceRequest request, + StreamObserver responseObserver) { CompletableFuture future = grpcForwardService.reportThreadStackTrace(Context.current(), request); future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) .exceptionally(e -> { @@ -264,7 +266,8 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic } @Override - public void reportMessageConsumptionResult(ReportMessageConsumptionResultRequest request, StreamObserver responseObserver) { + public void reportMessageConsumptionResult(ReportMessageConsumptionResultRequest request, + StreamObserver responseObserver) { CompletableFuture future = grpcForwardService.reportMessageConsumptionResult(Context.current(), request); future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) .exceptionally(e -> { @@ -277,7 +280,8 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic } @Override - public void notifyClientTermination(NotifyClientTerminationRequest request, StreamObserver responseObserver) { + public void notifyClientTermination(NotifyClientTerminationRequest request, + StreamObserver responseObserver) { CompletableFuture future = grpcForwardService.notifyClientTermination(Context.current(), request); future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) .exceptionally(e -> { @@ -290,7 +294,8 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic } @Override - public void changeInvisibleDuration(ChangeInvisibleDurationRequest request, StreamObserver responseObserver) { + public void changeInvisibleDuration(ChangeInvisibleDurationRequest request, + StreamObserver responseObserver) { CompletableFuture future = grpcForwardService.changeInvisibleDuration(Context.current(), request); future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) .exceptionally(e -> { diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcServer.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcServer.java index 6ebda585a4..8be76a95d3 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcServer.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/GrpcServer.java @@ -97,6 +97,7 @@ public class GrpcServer implements StartAndShutdown { .addService(messagingProcessor) .executor(this.executor); + // grpc interceptors, including acl, logging etc. if (ConfigurationManager.getProxyConfig().isEnableACL()) { List accessValidators = ServiceProvider.load(ServiceProvider.ACL_VALIDATOR_ID, AccessValidator.class); if (accessValidators.isEmpty()) { @@ -122,7 +123,7 @@ public class GrpcServer implements StartAndShutdown { this.grpcForwardService.start(); this.server.start(); - log.info("grpc server has started"); + log.info("grpc server start successfully."); } public void shutdown() { @@ -132,7 +133,7 @@ public class GrpcServer implements StartAndShutdown { this.grpcForwardService.shutdown(); - log.info("grpc server has stopped"); + log.info("grpc server shutdown successfully."); } catch (Exception e) { e.printStackTrace(); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/GrpcConverter.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/GrpcConverter.java index 7199069975..8b5265a049 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/GrpcConverter.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/GrpcConverter.java @@ -175,8 +175,7 @@ public class GrpcConverter { public static AckMessageRequestHeader buildAckMessageRequestHeader(AckMessageRequest request) { String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); - String receiptHandleStr = request.getReceiptHandle(); - ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr); + ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle()); AckMessageRequestHeader ackMessageRequestHeader = new AckMessageRequestHeader(); ackMessageRequestHeader.setConsumerGroup(groupName); @@ -191,8 +190,7 @@ public class GrpcConverter { DelayPolicy delayPolicy) { String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); - String receiptHandleStr = request.getReceiptHandle(); - ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr); + ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle()); ChangeInvisibleTimeRequestHeader changeInvisibleTimeRequestHeader = new ChangeInvisibleTimeRequestHeader(); changeInvisibleTimeRequestHeader.setConsumerGroup(groupName); @@ -209,8 +207,7 @@ public class GrpcConverter { ChangeInvisibleDurationRequest request) { String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); - String receiptHandleStr = request.getReceiptHandle(); - ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr); + ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle()); ChangeInvisibleTimeRequestHeader changeInvisibleTimeRequestHeader = new ChangeInvisibleTimeRequestHeader(); changeInvisibleTimeRequestHeader.setConsumerGroup(groupName); @@ -226,8 +223,7 @@ public class GrpcConverter { ForwardMessageToDeadLetterQueueRequest request) { String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); - String receiptHandleStr = request.getReceiptHandle(); - ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr); + ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle()); ConsumerSendMsgBackRequestHeader consumerSendMsgBackRequestHeader = new ConsumerSendMsgBackRequestHeader(); consumerSendMsgBackRequestHeader.setOffset(handle.getCommitLogOffset()); @@ -243,8 +239,7 @@ public class GrpcConverter { NackMessageRequest request) { String groupName = GrpcConverter.wrapResourceWithNamespace(request.getGroup()); String topicName = GrpcConverter.wrapResourceWithNamespace(request.getTopic()); - String receiptHandleStr = request.getReceiptHandle(); - ReceiptHandle handle = ReceiptHandle.decode(receiptHandleStr); + ReceiptHandle handle = ReceiptHandle.decode(request.getReceiptHandle()); ConsumerSendMsgBackRequestHeader consumerSendMsgBackRequestHeader = new ConsumerSendMsgBackRequestHeader(); consumerSendMsgBackRequestHeader.setOffset(handle.getCommitLogOffset()); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientService.java index 32e9b76b47..ca6f1e79b4 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ForwardClientService.java @@ -68,10 +68,12 @@ public class ForwardClientService extends BaseService { this.pollCommandResponseManager = pollCommandResponseManager; this.consumerManager = new ConsumerManager(new ConsumerIdsChangeListener() { - @Override public void handle(ConsumerGroupEvent event, String group, Object... args) { + @Override + public void handle(ConsumerGroupEvent event, String group, Object... args) { } - @Override public void shutdown() { + @Override + public void shutdown() { } }); this.producerManager = new ProducerManager(); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/connector/transaction/TransactionIdTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/connector/transaction/TransactionIdTest.java index 4dda6996c8..ab4cfb0201 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/connector/transaction/TransactionIdTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/connector/transaction/TransactionIdTest.java @@ -10,7 +10,7 @@ public class TransactionIdTest { @Test public void test() throws UnknownHostException { - TransactionId transactionId = TransactionId.genFromBrokerTransactionId( + TransactionId transactionId = TransactionId.genByBrokerTransactionId( RemotingHelper.string2SocketAddress("127.0.0.1:8080"), "71F99B78B6E261357FA259CCA6456118", 1234, 5678); @@ -24,7 +24,7 @@ public class TransactionIdTest { @Test public void testEmptyTransactionId() throws UnknownHostException { - TransactionId transactionId = TransactionId.genFromBrokerTransactionId( + TransactionId transactionId = TransactionId.genByBrokerTransactionId( RemotingHelper.string2SocketAddress("127.0.0.1:8080"), "", 1234, 5678); @@ -38,7 +38,7 @@ public class TransactionIdTest { @Test public void testNullTransactionId() throws UnknownHostException { - TransactionId transactionId = TransactionId.genFromBrokerTransactionId( + TransactionId transactionId = TransactionId.genByBrokerTransactionId( RemotingHelper.string2SocketAddress("127.0.0.1:8080"), null, 1234, 5678); diff --git a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcServiceTest.java b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcServiceTest.java index 6a252d2ca8..98b96c123a 100644 --- a/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcServiceTest.java +++ b/proxy/src/test/java/org/apache/rocketmq/proxy/grpc/service/LocalGrpcServiceTest.java @@ -381,7 +381,7 @@ public class LocalGrpcServiceTest extends InitConfigAndLoggerTest { .thenReturn(response); EndTransactionRequest request = EndTransactionRequest.newBuilder() .setMessageId("123") - .setTransactionId(TransactionId.genFromBrokerTransactionId( + .setTransactionId(TransactionId.genByBrokerTransactionId( new InetSocketAddress("0.0.0.0", 80), "123", 123, 123 ).getProxyTransactionId() ) 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 e54529f5e7..73f243b7b0 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 @@ -48,7 +48,7 @@ public class TransactionServiceTest extends BaseServiceTest { return null; }).when(channel).writeAndFlush(any()); - TransactionId transactionId = TransactionId.genFromBrokerTransactionId( + TransactionId transactionId = TransactionId.genByBrokerTransactionId( RemotingHelper.string2SocketAddress("127.0.0.1:8080"), "71F99B78B6E261357FA259CCA6456118", 1234, 5678); transactionService.checkTransactionState(new TransactionStateCheckRequest( @@ -69,7 +69,7 @@ public class TransactionServiceTest extends BaseServiceTest { public void testEndTransaction() throws Exception { AtomicReference headerRef = new AtomicReference<>(); AtomicReference brokerAddrRef = new AtomicReference<>(); - TransactionId transactionId = TransactionId.genFromBrokerTransactionId( + TransactionId transactionId = TransactionId.genByBrokerTransactionId( RemotingHelper.string2SocketAddress("127.0.0.1:8080"), "71F99B78B6E261357FA259CCA6456118", 1234, 5678); doAnswer(mock -> {