From 7f4b7bce2df26aa2497bd7ce97064b41a6f835b5 Mon Sep 17 00:00:00 2001 From: "Jixiang.jjx" Date: Thu, 17 Mar 2022 21:28:39 +0800 Subject: [PATCH] [ISSUE #3949] Do refactor work for passing check style. --- .../grpc/service/ClusterGrpcService.java | 54 +++++++++++++++---- .../grpc/service/cluster/ClientService.java | 8 +-- .../grpc/service/cluster/ProducerService.java | 17 +++--- .../service/cluster/PullMessageService.java | 51 +++++++++--------- .../cluster/ReceiveMessageService.java | 53 +++++++++--------- .../grpc/service/cluster/RouteService.java | 41 ++++++++------ 6 files changed, 138 insertions(+), 86 deletions(-) diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java index 8135f13f75..7479cabf46 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/ClusterGrpcService.java @@ -17,7 +17,40 @@ package org.apache.rocketmq.proxy.grpc.service; -import apache.rocketmq.v1.*; +import apache.rocketmq.v1.AckMessageRequest; +import apache.rocketmq.v1.AckMessageResponse; +import apache.rocketmq.v1.ChangeInvisibleDurationRequest; +import apache.rocketmq.v1.ChangeInvisibleDurationResponse; +import apache.rocketmq.v1.EndTransactionRequest; +import apache.rocketmq.v1.EndTransactionResponse; +import apache.rocketmq.v1.ForwardMessageToDeadLetterQueueRequest; +import apache.rocketmq.v1.ForwardMessageToDeadLetterQueueResponse; +import apache.rocketmq.v1.HealthCheckRequest; +import apache.rocketmq.v1.HealthCheckResponse; +import apache.rocketmq.v1.HeartbeatRequest; +import apache.rocketmq.v1.HeartbeatResponse; +import apache.rocketmq.v1.NackMessageRequest; +import apache.rocketmq.v1.NackMessageResponse; +import apache.rocketmq.v1.NotifyClientTerminationRequest; +import apache.rocketmq.v1.NotifyClientTerminationResponse; +import apache.rocketmq.v1.PollCommandRequest; +import apache.rocketmq.v1.PollCommandResponse; +import apache.rocketmq.v1.PullMessageRequest; +import apache.rocketmq.v1.PullMessageResponse; +import apache.rocketmq.v1.QueryAssignmentRequest; +import apache.rocketmq.v1.QueryAssignmentResponse; +import apache.rocketmq.v1.QueryOffsetRequest; +import apache.rocketmq.v1.QueryOffsetResponse; +import apache.rocketmq.v1.QueryRouteRequest; +import apache.rocketmq.v1.QueryRouteResponse; +import apache.rocketmq.v1.ReceiveMessageRequest; +import apache.rocketmq.v1.ReceiveMessageResponse; +import apache.rocketmq.v1.ReportMessageConsumptionResultRequest; +import apache.rocketmq.v1.ReportMessageConsumptionResultResponse; +import apache.rocketmq.v1.ReportThreadStackTraceRequest; +import apache.rocketmq.v1.ReportThreadStackTraceResponse; +import apache.rocketmq.v1.SendMessageRequest; +import apache.rocketmq.v1.SendMessageResponse; import com.google.rpc.Code; import io.grpc.Context; import org.apache.rocketmq.common.ThreadFactoryImpl; @@ -28,10 +61,13 @@ import org.apache.rocketmq.proxy.common.StartAndShutdown; import org.apache.rocketmq.proxy.connector.ConnectorManager; 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.common.Converter; import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; -import org.apache.rocketmq.proxy.grpc.service.cluster.*; +import org.apache.rocketmq.proxy.grpc.service.cluster.ClientService; +import org.apache.rocketmq.proxy.grpc.service.cluster.ProducerService; +import org.apache.rocketmq.proxy.grpc.service.cluster.PullMessageService; +import org.apache.rocketmq.proxy.grpc.service.cluster.ReceiveMessageService; +import org.apache.rocketmq.proxy.grpc.service.cluster.RouteService; +import org.apache.rocketmq.proxy.grpc.service.cluster.TransactionService; import org.slf4j.Logger; import org.slf4j.LoggerFactory; @@ -116,7 +152,7 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc @Override public CompletableFuture forwardMessageToDeadLetterQueue(Context ctx, - ForwardMessageToDeadLetterQueueRequest request) { + ForwardMessageToDeadLetterQueueRequest request) { return this.producerService.forwardMessageToDeadLetterQueue(ctx, request); } @@ -141,19 +177,19 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc } @Override public CompletableFuture reportThreadStackTrace(Context ctx, - ReportThreadStackTraceRequest request) { + ReportThreadStackTraceRequest request) { return null; } @Override public CompletableFuture reportMessageConsumptionResult(Context ctx, - ReportMessageConsumptionResultRequest request) { + ReportMessageConsumptionResultRequest request) { return null; } @Override public CompletableFuture notifyClientTermination(Context ctx, - NotifyClientTerminationRequest request) { + NotifyClientTerminationRequest request) { this.clientService.unregister(ctx, request, channelManager); return CompletableFuture.completedFuture(NotifyClientTerminationResponse.newBuilder() .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) @@ -161,7 +197,7 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc } @Override public CompletableFuture changeInvisibleDuration(Context ctx, - ChangeInvisibleDurationRequest request) { + ChangeInvisibleDurationRequest request) { return null; } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientService.java index 2eef9a38d7..95fa09acf3 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ClientService.java @@ -24,16 +24,12 @@ import apache.rocketmq.v1.PollCommandRequest; import apache.rocketmq.v1.PollCommandResponse; import apache.rocketmq.v1.Resource; import io.grpc.Context; -import java.util.concurrent.CompletableFuture; -import java.util.concurrent.ScheduledExecutorService; -import java.util.concurrent.TimeUnit; import org.apache.rocketmq.broker.client.ClientChannelInfo; import org.apache.rocketmq.broker.client.ConsumerManager; import org.apache.rocketmq.broker.client.ProducerManager; import org.apache.rocketmq.common.MQVersion; import org.apache.rocketmq.proxy.channel.ChannelManager; import org.apache.rocketmq.proxy.connector.ConnectorManager; -import org.apache.rocketmq.proxy.connector.transaction.TransactionHeartbeatRegisterService; import org.apache.rocketmq.proxy.grpc.adapter.channel.GrpcClientChannel; import org.apache.rocketmq.proxy.grpc.common.Converter; import org.apache.rocketmq.proxy.grpc.common.InterceptorConstants; @@ -41,6 +37,10 @@ import org.apache.rocketmq.remoting.protocol.LanguageCode; import org.slf4j.Logger; import org.slf4j.LoggerFactory; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.ScheduledExecutorService; +import java.util.concurrent.TimeUnit; + public class ClientService extends BaseService { private static final Logger log = LoggerFactory.getLogger(ClientService.class); 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 bcf949f872..1ee3d36c9d 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 @@ -23,7 +23,6 @@ import apache.rocketmq.v1.SendMessageRequest; import apache.rocketmq.v1.SendMessageResponse; import com.google.rpc.Code; import io.grpc.Context; -import java.util.concurrent.CompletableFuture; import org.apache.commons.lang3.StringUtils; import org.apache.commons.lang3.tuple.Pair; import org.apache.rocketmq.client.producer.SendResult; @@ -31,15 +30,17 @@ 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.route.SelectableMessageQueue; -import org.apache.rocketmq.proxy.common.utils.ProxyUtils; import org.apache.rocketmq.proxy.grpc.common.Converter; import org.apache.rocketmq.proxy.grpc.common.ProxyException; import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.common.ResponseHook; import org.apache.rocketmq.remoting.protocol.RemotingCommand; +import java.util.concurrent.CompletableFuture; + public class ProducerService extends BaseService { private volatile ProducerQueueSelector messageQueueSelector; @@ -149,14 +150,12 @@ public class ProducerService extends BaseService { ConsumerSendMsgBackRequestHeader requestHeader = this.convertToConsumerSendMsgBackRequestHeader(ctx, request); CompletableFuture resultFuture = this.connectorManager.getForwardProducer() .sendMessageBack(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); - resultFuture.thenAccept(result -> { - future.complete(ForwardMessageToDeadLetterQueueResponse.newBuilder() + resultFuture.thenAccept(result -> future.complete(ForwardMessageToDeadLetterQueueResponse.newBuilder() .setCommon(ResponseBuilder.buildCommon(result.getCode(), result.getRemark())) - .build()); - }).exceptionally(throwable -> { - future.completeExceptionally(throwable); - return null; - }); + .build())).exceptionally(throwable -> { + future.completeExceptionally(throwable); + return null; + }); } catch (Throwable t) { future.completeExceptionally(t); } 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 e406410ffa..587059b38b 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 @@ -26,10 +26,6 @@ import apache.rocketmq.v1.QueryOffsetResponse; import com.google.protobuf.util.Timestamps; import com.google.rpc.Code; import io.grpc.Context; -import java.util.ArrayList; -import java.util.List; -import java.util.concurrent.CompletableFuture; -import java.util.concurrent.TimeUnit; import org.apache.rocketmq.client.consumer.PullResult; import org.apache.rocketmq.client.consumer.PullStatus; import org.apache.rocketmq.common.message.MessageExt; @@ -45,6 +41,11 @@ import org.apache.rocketmq.proxy.grpc.common.ProxyException; import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.common.ResponseHook; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; + public class PullMessageService extends BaseService { private final DefaultForwardClient defaultForwardClient; @@ -82,15 +83,13 @@ public class PullMessageService extends BaseService { String brokerAddr = this.getBrokerAddr(ctx, brokerName); offsetFuture = this.defaultForwardClient.searchOffset(brokerAddr, topic, queueId, timestamp, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); } - offsetFuture.thenAccept(result -> { - future.complete(QueryOffsetResponse.newBuilder() - .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) - .setOffset(result) - .build()); - }).exceptionally(throwable -> { - future.completeExceptionally(throwable); - return null; - }); + offsetFuture.thenAccept(result -> future.complete(QueryOffsetResponse.newBuilder() + .setCommon(ResponseBuilder.buildCommon(Code.OK, Code.OK.name())) + .setOffset(result) + .build())).exceptionally(throwable -> { + future.completeExceptionally(throwable); + return null; + }); } catch (Throwable t) { future.completeExceptionally(t); } @@ -104,6 +103,7 @@ public class PullMessageService extends BaseService { pullMessageHook.beforeResponse(request, response, throwable); } }); + try { PullMessageRequestHeader requestHeader = this.convertToPullMessageRequestHeader(ctx, request); @@ -111,17 +111,20 @@ public class PullMessageService extends BaseService { String brokerAddr = this.getBrokerAddr(ctx, brokerName); CompletableFuture pullResultFuture = this.connectorManager.getForwardReadConsumer() - .pullMessage(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); - pullResultFuture.thenAccept(pullResult -> { - try { - future.complete(convertToPullMessageResponse(ctx, request, pullResult)); - } catch (Throwable throwable) { - future.completeExceptionally(throwable); - } - }).exceptionally(throwable -> { - future.completeExceptionally(throwable); - return null; - }); + .pullMessage(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); + pullResultFuture + .thenAccept(pullResult -> { + try { + future.complete(convertToPullMessageResponse(ctx, request, pullResult)); + } catch (Throwable throwable) { + future.completeExceptionally(throwable); + } + }) + .exceptionally(throwable -> { + future.completeExceptionally(throwable); + return null; + }); + } catch (Throwable t) { future.completeExceptionally(t); } diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ReceiveMessageService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ReceiveMessageService.java index 0019656611..ffd91a14c3 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ReceiveMessageService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/service/cluster/ReceiveMessageService.java @@ -25,10 +25,6 @@ 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 java.util.concurrent.TimeUnit; import org.apache.rocketmq.client.consumer.AckResult; import org.apache.rocketmq.client.consumer.AckStatus; import org.apache.rocketmq.client.consumer.PopResult; @@ -48,6 +44,11 @@ import org.apache.rocketmq.proxy.grpc.common.Converter; import org.apache.rocketmq.proxy.grpc.common.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.common.ResponseHook; +import java.util.ArrayList; +import java.util.List; +import java.util.concurrent.CompletableFuture; +import java.util.concurrent.TimeUnit; + public class ReceiveMessageService extends BaseService { private final ForwardReadConsumer readConsumer; @@ -147,16 +148,18 @@ public class ReceiveMessageService extends BaseService { AckMessageRequestHeader requestHeader = this.convertToAckMessageRequestHeader(ctx, request); CompletableFuture ackResultFuture = this.writeConsumer.ackMessage(brokerAddr, requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); - ackResultFuture.thenAccept(result -> { - try { - future.complete(convertToAckMessageResponse(ctx, request, result)); - } catch (Throwable throwable) { - future.completeExceptionally(throwable); - } - }).exceptionally(throwable -> { - future.completeExceptionally(throwable); - return null; - }); + ackResultFuture + .thenAccept(result -> { + try { + future.complete(convertToAckMessageResponse(ctx, request, result)); + } catch (Throwable throwable) { + future.completeExceptionally(throwable); + } + }) + .exceptionally(throwable -> { + future.completeExceptionally(throwable); + return null; + }); } catch (Throwable t) { future.completeExceptionally(t); } @@ -192,16 +195,18 @@ public class ReceiveMessageService extends BaseService { ChangeInvisibleTimeRequestHeader requestHeader = this.convertToChangeInvisibleTimeRequestHeader(ctx, request); CompletableFuture resultFuture = this.writeConsumer.changeInvisibleTimeAsync(brokerAddr, receiptHandle.getBrokerName(), requestHeader, ProxyUtils.DEFAULT_MQ_CLIENT_TIMEOUT); - resultFuture.thenAccept(result -> { - try { - future.complete(convertToNackMessageResponse(ctx, request, result)); - } catch (Throwable throwable) { - future.completeExceptionally(throwable); - } - }).exceptionally(throwable -> { - future.completeExceptionally(throwable); - return null; - }); + resultFuture + .thenAccept(result -> { + try { + future.complete(convertToNackMessageResponse(ctx, request, result)); + } catch (Throwable throwable) { + future.completeExceptionally(throwable); + } + }) + .exceptionally(throwable -> { + future.completeExceptionally(throwable); + return null; + }); } catch (Throwable t) { future.completeExceptionally(t); } 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 ea8ad005fe..4adde9ea49 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 @@ -16,7 +16,16 @@ */ package org.apache.rocketmq.proxy.grpc.service.cluster; -import apache.rocketmq.v1.*; +import apache.rocketmq.v1.Assignment; +import apache.rocketmq.v1.Broker; +import apache.rocketmq.v1.Endpoints; +import apache.rocketmq.v1.Partition; +import apache.rocketmq.v1.Permission; +import apache.rocketmq.v1.QueryAssignmentRequest; +import apache.rocketmq.v1.QueryAssignmentResponse; +import apache.rocketmq.v1.QueryRouteRequest; +import apache.rocketmq.v1.QueryRouteResponse; +import apache.rocketmq.v1.Resource; import com.google.rpc.Code; import io.grpc.Context; import org.apache.rocketmq.common.constant.PermName; @@ -140,35 +149,35 @@ public class RouteService extends BaseService { r = queueData.getReadQueueNums(); } - // r here means readOnly queue nums, w means writeOnly queue nums, while rw means readable and writable queue nums. + // r here means readOnly queue nums, w means writeOnly queue nums, while rw means both readable and writable queue nums. int queueIdIndex = 0; - for(int i = 0; i < r; i++){ - Partition partition = buildPartition(broker, topic, queueIdIndex++, Permission.READ); + for (int i = 0; i < r; i++) { + Partition partition = Partition.newBuilder().setBroker(broker).setTopic(topic) + .setId(queueIdIndex++) + .setPermission(Permission.READ) + .build(); partitionList.add(partition); } - for(int i = 0; i < w; i++){ - Partition partition = buildPartition(broker, topic, queueIdIndex++, Permission.WRITE); + for (int i = 0; i < w; i++) { + Partition partition = Partition.newBuilder().setBroker(broker).setTopic(topic) + .setId(queueIdIndex++) + .setPermission(Permission.WRITE) + .build(); partitionList.add(partition); } for (int i = 0; i < rw; i++) { - Partition partition = buildPartition(broker, topic, queueIdIndex++, Permission.READ_WRITE); + Partition partition = Partition.newBuilder().setBroker(broker).setTopic(topic) + .setId(queueIdIndex++) + .setPermission(Permission.READ_WRITE) + .build(); partitionList.add(partition); } return partitionList; } - private static Partition buildPartition(Broker broker, Resource topic, int queueId, Permission perm) { - Partition.Builder builder = Partition.newBuilder() - .setBroker(broker) - .setTopic(topic) - .setId(queueId); - builder.setPermission(perm); - return builder.build(); - } - public CompletableFuture queryAssignment(Context ctx, QueryAssignmentRequest request) { CompletableFuture future = new CompletableFuture<>(); future.whenComplete((response, throwable) -> {