diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessagingProcessor.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessagingProcessor.java index 6317cddae3..586061b3e4 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessagingProcessor.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/GrpcMessagingProcessor.java @@ -117,15 +117,7 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic @Override public void receiveMessage(ReceiveMessageRequest request, StreamObserver responseObserver) { - CompletableFuture> future = grpcForwardService.receiveMessage(Context.current(), request); - future.thenAccept(response -> ResponseWriter.write(responseObserver, response)) - .exceptionally(e -> { - ResponseWriter.write( - responseObserver, - ReceiveMessageResponse.newBuilder().setStatus(convertExceptionToStatus(e)).build() - ); - return null; - }); + grpcForwardService.receiveMessage(Context.current(), request, responseObserver); } @Override diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java index 7dfd982981..9c4d663d67 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/ClusterGrpcService.java @@ -42,8 +42,6 @@ import apache.rocketmq.v2.SendMessageResponse; import apache.rocketmq.v2.TelemetryCommand; import io.grpc.Context; import io.grpc.stub.StreamObserver; -import java.util.Iterator; -import java.util.List; import java.util.concurrent.CompletableFuture; import java.util.concurrent.Executors; import java.util.concurrent.ScheduledExecutorService; @@ -123,8 +121,9 @@ public class ClusterGrpcService extends AbstractStartAndShutdown implements Grpc } @Override - public CompletableFuture> receiveMessage(Context ctx, ReceiveMessageRequest request) { - return consumerService.receiveMessage(ctx, request).thenApply(List::iterator); + public void receiveMessage(Context ctx, ReceiveMessageRequest request, + StreamObserver responseObserver) { + consumerService.receiveMessage(ctx, request, responseObserver); } @Override diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcForwardService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcForwardService.java index b1c984defe..48d2f2504f 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcForwardService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/GrpcForwardService.java @@ -42,7 +42,6 @@ import apache.rocketmq.v2.SendMessageResponse; import apache.rocketmq.v2.TelemetryCommand; import io.grpc.Context; import io.grpc.stub.StreamObserver; -import java.util.Iterator; import java.util.concurrent.CompletableFuture; import org.apache.rocketmq.proxy.common.StartAndShutdown; @@ -56,7 +55,7 @@ public interface GrpcForwardService extends StartAndShutdown { CompletableFuture queryAssignment(Context ctx, QueryAssignmentRequest request); - CompletableFuture> receiveMessage(Context ctx, ReceiveMessageRequest request); + void receiveMessage(Context ctx, ReceiveMessageRequest request, StreamObserver responseObserver); CompletableFuture nackMessage(Context ctx, NackMessageRequest request); diff --git a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java index dcbbf600c4..fb8a277353 100644 --- a/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java +++ b/proxy/src/main/java/org/apache/rocketmq/proxy/grpc/v2/service/cluster/ConsumerService.java @@ -33,6 +33,7 @@ import apache.rocketmq.v2.Resource; import apache.rocketmq.v2.RetryPolicy; import apache.rocketmq.v2.Settings; import io.grpc.Context; +import io.grpc.stub.StreamObserver; import java.util.ArrayList; import java.util.Collections; import java.util.List; @@ -59,6 +60,7 @@ import org.apache.rocketmq.proxy.grpc.v2.adapter.GrpcConverter; import org.apache.rocketmq.proxy.grpc.v2.adapter.ProxyException; import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseBuilder; import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseHook; +import org.apache.rocketmq.proxy.grpc.v2.adapter.ResponseWriter; import org.apache.rocketmq.proxy.grpc.v2.service.GrpcClientManager; import org.apache.rocketmq.remoting.protocol.RemotingCommand; @@ -90,7 +92,20 @@ public class ConsumerService extends BaseService { this.grpcClientManager = grpcClientManager; } - public CompletableFuture> receiveMessage(Context ctx, ReceiveMessageRequest request) { + public void receiveMessage(Context ctx, ReceiveMessageRequest request, + StreamObserver responseObserver) { + this.receiveMessage(ctx, request) + .thenAccept(responses -> ResponseWriter.write(responseObserver, responses.iterator())) + .exceptionally(e -> { + ResponseWriter.write( + responseObserver, + ReceiveMessageResponse.newBuilder().setStatus(ResponseBuilder.buildStatus(e)).build() + ); + return null; + }); + } + + protected CompletableFuture> receiveMessage(Context ctx, ReceiveMessageRequest request) { CompletableFuture> future = new CompletableFuture<>(); try {