mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] Do refactor work for passing check style.
This commit is contained in:
@@ -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<ForwardMessageToDeadLetterQueueResponse> 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<ReportThreadStackTraceResponse> reportThreadStackTrace(Context ctx,
|
||||
ReportThreadStackTraceRequest request) {
|
||||
ReportThreadStackTraceRequest request) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public CompletableFuture<ReportMessageConsumptionResultResponse> reportMessageConsumptionResult(Context ctx,
|
||||
ReportMessageConsumptionResultRequest request) {
|
||||
ReportMessageConsumptionResultRequest request) {
|
||||
return null;
|
||||
}
|
||||
|
||||
@Override
|
||||
public CompletableFuture<NotifyClientTerminationResponse> 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<ChangeInvisibleDurationResponse> changeInvisibleDuration(Context ctx,
|
||||
ChangeInvisibleDurationRequest request) {
|
||||
ChangeInvisibleDurationRequest request) {
|
||||
return null;
|
||||
}
|
||||
|
||||
|
||||
+4
-4
@@ -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);
|
||||
|
||||
+8
-9
@@ -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<RemotingCommand> 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);
|
||||
}
|
||||
|
||||
+27
-24
@@ -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<PullResult> 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);
|
||||
}
|
||||
|
||||
+29
-24
@@ -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<AckResult> 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<AckResult> 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);
|
||||
}
|
||||
|
||||
+25
-16
@@ -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<QueryAssignmentResponse> queryAssignment(Context ctx, QueryAssignmentRequest request) {
|
||||
CompletableFuture<QueryAssignmentResponse> future = new CompletableFuture<>();
|
||||
future.whenComplete((response, throwable) -> {
|
||||
|
||||
Reference in New Issue
Block a user