From 95f99845713ac4ab98201c4e83052411a23736f6 Mon Sep 17 00:00:00 2001 From: caoyanan666 <55247691+caoyanan666@users.noreply.github.com> Date: Fri, 31 May 2024 16:26:59 +0800 Subject: [PATCH] fix triple reactor call throws "Too many response for unary method" exception (#14261) Co-authored-by: caoyanan --- .../apache/dubbo/reactive/calls/ReactorServerCalls.java | 8 +++----- 1 file changed, 3 insertions(+), 5 deletions(-) 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 24218a1b08..8cf1ef3ed8 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 @@ -48,13 +48,11 @@ public final class ReactorServerCalls { public static void oneToOne(T request, StreamObserver responseObserver, Function, Mono> func) { try { func.apply(Mono.just(request)) + .switchIfEmpty(Mono.error(TriRpcStatus.NOT_FOUND.asException())) .subscribe( - res -> { - responseObserver.onNext(res); - responseObserver.onCompleted(); - }, + responseObserver::onNext, throwable -> doOnResponseHasException(throwable, responseObserver), - () -> doOnResponseHasException(TriRpcStatus.NOT_FOUND.asException(), responseObserver)); + responseObserver::onCompleted); } catch (Throwable throwable) { doOnResponseHasException(throwable, responseObserver); }