From 9cf17ebe495739801d1a386b74fd2f4c9156e065 Mon Sep 17 00:00:00 2001 From: "Jixiang.jjx" Date: Mon, 28 Mar 2022 16:45:17 +0800 Subject: [PATCH] [ISSUE #3949] Do some refactoring work. --- .../rocketmq/client/impl/MQClientAPIExt.java | 29 ++++++++++--------- .../proxy/connector/ForwardProducer.java | 1 - .../connector/transaction/TransactionId.java | 9 ++++-- .../grpc/adapter/channel/ChannelType.java | 2 +- .../adapter/channel/GrpcClientChannel.java | 3 +- .../service/cluster/PullMessageService.java | 1 - 6 files changed, 24 insertions(+), 21 deletions(-) diff --git a/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExt.java b/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExt.java index c440cd6c0b..19cbb03970 100644 --- a/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExt.java +++ b/client/src/main/java/org/apache/rocketmq/client/impl/MQClientAPIExt.java @@ -180,9 +180,10 @@ public class MQClientAPIExt { ConsumerSendMsgBackRequestHeader requestHeader, long timeoutMillis ) { + RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.CONSUMER_SEND_MSG_BACK, requestHeader); + CompletableFuture future = new CompletableFuture<>(); try { - RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.CONSUMER_SEND_MSG_BACK, requestHeader); this.getRemotingClient().invokeAsync(brokerAddr, request, timeoutMillis, responseFuture -> { RemotingCommand response = responseFuture.getResponseCommand(); if (response != null) { @@ -357,19 +358,19 @@ public class MQClientAPIExt { } public CompletableFuture getMaxOffset(String brokerAddr, String topic, int queueId, long timeoutMillis) { + GetMaxOffsetRequestHeader requestHeader = new GetMaxOffsetRequestHeader(); + requestHeader.setTopic(topic); + requestHeader.setQueueId(queueId); + RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.GET_MAX_OFFSET, requestHeader); + CompletableFuture future = new CompletableFuture<>(); try { - GetMaxOffsetRequestHeader requestHeader = new GetMaxOffsetRequestHeader(); - requestHeader.setTopic(topic); - requestHeader.setQueueId(queueId); - RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.GET_MAX_OFFSET, requestHeader); this.getRemotingClient().invokeAsync(brokerAddr, request, timeoutMillis, responseFuture -> { RemotingCommand response = responseFuture.getResponseCommand(); if (response != null) { if (ResponseCode.SUCCESS == response.getCode()) { try { - GetMaxOffsetResponseHeader responseHeader = - (GetMaxOffsetResponseHeader) response.decodeCommandCustomHeader(GetMaxOffsetResponseHeader.class); + GetMaxOffsetResponseHeader responseHeader = response.decodeCommandCustomHeader(GetMaxOffsetResponseHeader.class); future.complete(responseHeader.getOffset()); } catch (Throwable t) { future.completeExceptionally(t); @@ -387,20 +388,20 @@ public class MQClientAPIExt { } public CompletableFuture searchOffset(String brokerAddr, String topic, int queueId , long timestamp, long timeoutMillis) { + SearchOffsetRequestHeader requestHeader = new SearchOffsetRequestHeader(); + requestHeader.setTopic(topic); + requestHeader.setQueueId(queueId); + requestHeader.setTimestamp(timestamp); + RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.SEARCH_OFFSET_BY_TIMESTAMP, requestHeader); + CompletableFuture future = new CompletableFuture<>(); try { - SearchOffsetRequestHeader requestHeader = new SearchOffsetRequestHeader(); - requestHeader.setTopic(topic); - requestHeader.setQueueId(queueId); - requestHeader.setTimestamp(timestamp); - RemotingCommand request = RemotingCommand.createRequestCommand(RequestCode.SEARCH_OFFSET_BY_TIMESTAMP, requestHeader); this.getRemotingClient().invokeAsync(brokerAddr, request, timeoutMillis, responseFuture -> { RemotingCommand response = responseFuture.getResponseCommand(); if (response != null) { if (response.getCode() == ResponseCode.SUCCESS) { try { - SearchOffsetResponseHeader responseHeader = - (SearchOffsetResponseHeader) response.decodeCommandCustomHeader(SearchOffsetResponseHeader.class); + SearchOffsetResponseHeader responseHeader = response.decodeCommandCustomHeader(SearchOffsetResponseHeader.class); future.complete(responseHeader.getOffset()); } catch (Throwable t) { future.completeExceptionally(t); 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 f31fe95c32..a6f8a0a6e7 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 @@ -32,7 +32,6 @@ import org.apache.rocketmq.proxy.connector.transaction.TransactionId; import org.apache.rocketmq.remoting.protocol.RemotingCommand; public class ForwardProducer extends AbstractForwardClient { - private static final String PID_PREFIX = "PID_RMQ_PROXY_PUBLISH_MESSAGE_"; public ForwardProducer(ForwardClientManager clientFactory) { 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 387bb5800d..e680c7ad90 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 @@ -42,8 +42,13 @@ public class TransactionId { private long tranStateTableOffset; private String proxyTransactionId; - public TransactionId(SocketAddress brokerAddr, String brokerTransactionId, long commitLogOffset, - long tranStateTableOffset, String proxyTransactionId) { + public TransactionId( + SocketAddress brokerAddr, + String brokerTransactionId, + long commitLogOffset, + long tranStateTableOffset, + String proxyTransactionId + ) { this.brokerAddr = brokerAddr; this.brokerTransactionId = brokerTransactionId; this.commitLogOffset = commitLogOffset; diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/ChannelType.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/ChannelType.java index 5a883b6a35..da8d30de8d 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/ChannelType.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/ChannelType.java @@ -23,7 +23,7 @@ public enum ChannelType { */ LOCAL, /** - * The channel sync from other proxy + * The channel synced from other proxy */ REMOTE } \ No newline at end of file diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannel.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannel.java index 6a3fcb045d..e389b02713 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannel.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/adapter/channel/GrpcClientChannel.java @@ -120,8 +120,7 @@ public class GrpcClientChannel extends SimpleChannel { try { switch (command.getCode()) { case RequestCode.CHECK_TRANSACTION_STATE: { - final CheckTransactionStateRequestHeader requestHeader = - (CheckTransactionStateRequestHeader) command.decodeCommandCustomHeader(CheckTransactionStateRequestHeader.class); + final CheckTransactionStateRequestHeader requestHeader = command.decodeCommandCustomHeader(CheckTransactionStateRequestHeader.class); MessageExt messageExt = MessageDecoder.decode(ByteBuffer.wrap(command.getBody()), true, false, false); future.complete(PollCommandResponse.newBuilder() .setRecoverOrphanedTransactionCommand(RecoverOrphanedTransactionCommand.newBuilder() 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 b3fc812089..b5b0545e19 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 @@ -41,7 +41,6 @@ import org.apache.rocketmq.proxy.grpc.adapter.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.adapter.ResponseHook; public class PullMessageService extends BaseService { - private final DefaultForwardClient forwardClient; private final ForwardReadConsumer readConsumer;