mirror of
https://github.com/apache/rocketmq.git
synced 2026-09-24 16:04:00 +08:00
[ISSUE #3949] v2 support
This commit is contained in:
@@ -117,15 +117,7 @@ public class GrpcMessagingProcessor extends MessagingServiceGrpc.MessagingServic
|
||||
|
||||
@Override
|
||||
public void receiveMessage(ReceiveMessageRequest request, StreamObserver<ReceiveMessageResponse> responseObserver) {
|
||||
CompletableFuture<Iterator<ReceiveMessageResponse>> 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
|
||||
|
||||
+3
-4
@@ -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<Iterator<ReceiveMessageResponse>> receiveMessage(Context ctx, ReceiveMessageRequest request) {
|
||||
return consumerService.receiveMessage(ctx, request).thenApply(List::iterator);
|
||||
public void receiveMessage(Context ctx, ReceiveMessageRequest request,
|
||||
StreamObserver<ReceiveMessageResponse> responseObserver) {
|
||||
consumerService.receiveMessage(ctx, request, responseObserver);
|
||||
}
|
||||
|
||||
@Override
|
||||
|
||||
+1
-2
@@ -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<QueryAssignmentResponse> queryAssignment(Context ctx, QueryAssignmentRequest request);
|
||||
|
||||
CompletableFuture<Iterator<ReceiveMessageResponse>> receiveMessage(Context ctx, ReceiveMessageRequest request);
|
||||
void receiveMessage(Context ctx, ReceiveMessageRequest request, StreamObserver<ReceiveMessageResponse> responseObserver);
|
||||
|
||||
CompletableFuture<NackMessageResponse> nackMessage(Context ctx, NackMessageRequest request);
|
||||
|
||||
|
||||
+16
-1
@@ -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<List<ReceiveMessageResponse>> receiveMessage(Context ctx, ReceiveMessageRequest request) {
|
||||
public void receiveMessage(Context ctx, ReceiveMessageRequest request,
|
||||
StreamObserver<ReceiveMessageResponse> 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<List<ReceiveMessageResponse>> receiveMessage(Context ctx, ReceiveMessageRequest request) {
|
||||
CompletableFuture<List<ReceiveMessageResponse>> future = new CompletableFuture<>();
|
||||
|
||||
try {
|
||||
|
||||
Reference in New Issue
Block a user