From 67440ebe27f5bc6cd43642467e023ee200b5dfc6 Mon Sep 17 00:00:00 2001 From: "kaiyi.lk" Date: Wed, 16 Mar 2022 16:57:13 +0800 Subject: [PATCH] [ISSUE #3949] pollCommand for cluster mode --- .../proxy/grpc/service/ClusterGrpcService.java | 16 ++++++++++++++-- 1 file changed, 14 insertions(+), 2 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 e2f0caf569..0769c6be09 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 @@ -31,6 +31,7 @@ import apache.rocketmq.v1.HeartbeatRequest; import apache.rocketmq.v1.HeartbeatResponse; import apache.rocketmq.v1.NackMessageRequest; import apache.rocketmq.v1.NackMessageResponse; +import apache.rocketmq.v1.NoopCommand; import apache.rocketmq.v1.NotifyClientTerminationRequest; import apache.rocketmq.v1.NotifyClientTerminationResponse; import apache.rocketmq.v1.PollCommandRequest; @@ -168,17 +169,28 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc public CompletableFuture pollCommand(Context ctx, PollCommandRequest request) { CompletableFuture future = new CompletableFuture<>(); String clientId = request.getClientId(); + PollCommandResponse noopCommandResponse = PollCommandResponse.newBuilder().setNoopCommand(NoopCommand.newBuilder().build()).build(); switch (request.getGroupCase()) { case PRODUCER_GROUP: Resource producerGroup = request.getProducerGroup(); String producerGroupName = Converter.getResourceNameWithNamespace(producerGroup); - GrpcClientChannel.addClientObserver(this.channelManager, producerGroupName, clientId, future); + GrpcClientChannel producerChannel = GrpcClientChannel.getChannel(this.channelManager, producerGroupName, clientId); + if (producerChannel == null) { + future.complete(noopCommandResponse); + } else { + producerChannel.addClientObserver(future); + } break; case CONSUMER_GROUP: Resource consumerGroup = request.getConsumerGroup(); String consumerGroupName = Converter.getResourceNameWithNamespace(consumerGroup); - GrpcClientChannel.addClientObserver(this.channelManager, consumerGroupName, clientId, future); + GrpcClientChannel consumerChannel = GrpcClientChannel.getChannel(this.channelManager, consumerGroupName, clientId); + if (consumerChannel == null) { + future.complete(noopCommandResponse); + } else { + consumerChannel.addClientObserver(future); + } break; default: break;