diff --git a/dubbo-plugin/dubbo-reactive/src/main/java/org/apache/dubbo/reactive/AbstractTripleReactorSubscriber.java b/dubbo-plugin/dubbo-reactive/src/main/java/org/apache/dubbo/reactive/AbstractTripleReactorSubscriber.java index 86db857582..b3d6fa058e 100644 --- a/dubbo-plugin/dubbo-reactive/src/main/java/org/apache/dubbo/reactive/AbstractTripleReactorSubscriber.java +++ b/dubbo-plugin/dubbo-reactive/src/main/java/org/apache/dubbo/reactive/AbstractTripleReactorSubscriber.java @@ -53,7 +53,7 @@ public abstract class AbstractTripleReactorSubscriber implements Subscriber extends AbstractTripleReactorSubscriber { + /** + * The execution future of the current task, in order to be returned to stubInvoker + */ + private final CompletableFuture> executionFuture = new CompletableFuture<>(); + /** + * The result elements collected by the current task. + * This class is a flux subscriber, which usually means there will be multiple elements, so it is declared as a list type. + */ + private final List collectedData = new ArrayList<>(); + + public ServerTripleReactorSubscriber() {} + + public ServerTripleReactorSubscriber(CallStreamObserver streamObserver) { + this.downstream = streamObserver; + } + @Override public void subscribe(CallStreamObserver downstream) { super.subscribe(downstream); @@ -40,4 +60,26 @@ public class ServerTripleReactorSubscriber extends AbstractTripleReactorSubsc context.addListener(ctx -> super.cancel()); } } + + @Override + public void onNext(T t) { + super.onNext(t); + collectedData.add(t); + } + + @Override + public void onError(Throwable throwable) { + super.onError(throwable); + executionFuture.completeExceptionally(throwable); + } + + @Override + public void onComplete() { + super.onComplete(); + executionFuture.complete(this.collectedData); + } + + public CompletableFuture> getExecutionFuture() { + return executionFuture; + } } diff --git a/dubbo-plugin/dubbo-reactive/src/main/java/org/apache/dubbo/reactive/calls/ReactorServerCalls.java b/dubbo-plugin/dubbo-reactive/src/main/java/org/apache/dubbo/reactive/calls/ReactorServerCalls.java index f6a39c944c..24218a1b08 100644 --- a/dubbo-plugin/dubbo-reactive/src/main/java/org/apache/dubbo/reactive/calls/ReactorServerCalls.java +++ b/dubbo-plugin/dubbo-reactive/src/main/java/org/apache/dubbo/reactive/calls/ReactorServerCalls.java @@ -24,6 +24,8 @@ import org.apache.dubbo.rpc.TriRpcStatus; import org.apache.dubbo.rpc.protocol.tri.observer.CallStreamObserver; import org.apache.dubbo.rpc.protocol.tri.observer.ServerCallToObserverAdapter; +import java.util.List; +import java.util.concurrent.CompletableFuture; import java.util.function.Function; import reactor.core.publisher.Flux; @@ -65,14 +67,21 @@ public final class ReactorServerCalls { * @param responseObserver response StreamObserver * @param func service implementation */ - public static void oneToMany( + public static CompletableFuture> oneToMany( T request, StreamObserver responseObserver, Function, Flux> func) { try { + ServerCallToObserverAdapter serverCallToObserverAdapter = + (ServerCallToObserverAdapter) responseObserver; Flux response = func.apply(Mono.just(request)); - ServerTripleReactorSubscriber subscriber = response.subscribeWith(new ServerTripleReactorSubscriber<>()); - subscriber.subscribe((ServerCallToObserverAdapter) responseObserver); + ServerTripleReactorSubscriber reactorSubscriber = + new ServerTripleReactorSubscriber<>(serverCallToObserverAdapter); + response.subscribeWith(reactorSubscriber).subscribe(serverCallToObserverAdapter); + return reactorSubscriber.getExecutionFuture(); } catch (Throwable throwable) { - responseObserver.onError(throwable); + doOnResponseHasException(throwable, responseObserver); + CompletableFuture> future = new CompletableFuture<>(); + future.completeExceptionally(throwable); + return future; } } diff --git a/dubbo-plugin/dubbo-reactive/src/main/java/org/apache/dubbo/reactive/handler/OneToManyMethodHandler.java b/dubbo-plugin/dubbo-reactive/src/main/java/org/apache/dubbo/reactive/handler/OneToManyMethodHandler.java index b5b0534fff..fc2e4df9f0 100644 --- a/dubbo-plugin/dubbo-reactive/src/main/java/org/apache/dubbo/reactive/handler/OneToManyMethodHandler.java +++ b/dubbo-plugin/dubbo-reactive/src/main/java/org/apache/dubbo/reactive/handler/OneToManyMethodHandler.java @@ -42,7 +42,6 @@ public class OneToManyMethodHandler implements StubMethodHandler { public CompletableFuture invoke(Object[] arguments) { T request = (T) arguments[0]; StreamObserver responseObserver = (StreamObserver) arguments[1]; - ReactorServerCalls.oneToMany(request, responseObserver, func); - return CompletableFuture.completedFuture(null); + return ReactorServerCalls.oneToMany(request, responseObserver, func); } }