fix triple reactor call throws "Too many response for unary method" exception (#14261)
Co-authored-by: caoyanan <caoyanan@growingio.com>
This commit is contained in:
parent
d35d7eb30c
commit
95f9984571
|
|
@ -48,13 +48,11 @@ public final class ReactorServerCalls {
|
||||||
public static <T, R> void oneToOne(T request, StreamObserver<R> responseObserver, Function<Mono<T>, Mono<R>> func) {
|
public static <T, R> void oneToOne(T request, StreamObserver<R> responseObserver, Function<Mono<T>, Mono<R>> func) {
|
||||||
try {
|
try {
|
||||||
func.apply(Mono.just(request))
|
func.apply(Mono.just(request))
|
||||||
|
.switchIfEmpty(Mono.error(TriRpcStatus.NOT_FOUND.asException()))
|
||||||
.subscribe(
|
.subscribe(
|
||||||
res -> {
|
responseObserver::onNext,
|
||||||
responseObserver.onNext(res);
|
|
||||||
responseObserver.onCompleted();
|
|
||||||
},
|
|
||||||
throwable -> doOnResponseHasException(throwable, responseObserver),
|
throwable -> doOnResponseHasException(throwable, responseObserver),
|
||||||
() -> doOnResponseHasException(TriRpcStatus.NOT_FOUND.asException(), responseObserver));
|
responseObserver::onCompleted);
|
||||||
} catch (Throwable throwable) {
|
} catch (Throwable throwable) {
|
||||||
doOnResponseHasException(throwable, responseObserver);
|
doOnResponseHasException(throwable, responseObserver);
|
||||||
}
|
}
|
||||||
|
|
|
||||||
Loading…
Reference in New Issue