From 84fb8550c776aaab48e943f24a1481041d8528f1 Mon Sep 17 00:00:00 2001 From: rodbate Date: Wed, 26 May 2021 17:50:55 +0800 Subject: [PATCH] Fix blocking discovery client. --- .../core/DiscoveryClientServiceInstanceListSupplier.java | 7 ++++--- .../DiscoveryClientServiceInstanceListSupplierTests.java | 1 + 2 files changed, 5 insertions(+), 3 deletions(-) diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DiscoveryClientServiceInstanceListSupplier.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DiscoveryClientServiceInstanceListSupplier.java index 49bb4058..42113416 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DiscoveryClientServiceInstanceListSupplier.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DiscoveryClientServiceInstanceListSupplier.java @@ -23,6 +23,7 @@ import java.util.List; import org.apache.commons.logging.Log; import org.apache.commons.logging.LogFactory; import reactor.core.publisher.Flux; +import reactor.core.publisher.Mono; import reactor.core.scheduler.Schedulers; import org.springframework.boot.convert.DurationStyle; @@ -39,6 +40,7 @@ import static org.springframework.cloud.loadbalancer.support.LoadBalancerClientF * @author Spencer Gibb * @author Olga Maciaszek-Sharma * @author Tim Ysewyn + * @author Rod Catter * @since 2.2.0 */ public class DiscoveryClientServiceInstanceListSupplier @@ -63,12 +65,11 @@ public class DiscoveryClientServiceInstanceListSupplier this.serviceId = environment.getProperty(PROPERTY_NAME); resolveTimeout(environment); this.serviceInstances = Flux - .defer(() -> Flux.just(delegate.getInstances(serviceId))) - .subscribeOn(Schedulers.boundedElastic()) + .defer(() -> Mono.fromCallable(() -> delegate.getInstances(serviceId))) .timeout(timeout, Flux.defer(() -> { logTimeout(); return Flux.just(new ArrayList<>()); - })).onErrorResume(error -> { + }), Schedulers.boundedElastic()).onErrorResume(error -> { logException(error); return Flux.just(new ArrayList<>()); }); diff --git a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/DiscoveryClientServiceInstanceListSupplierTests.java b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/DiscoveryClientServiceInstanceListSupplierTests.java index 980f7365..9c0ee468 100644 --- a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/DiscoveryClientServiceInstanceListSupplierTests.java +++ b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/DiscoveryClientServiceInstanceListSupplierTests.java @@ -39,6 +39,7 @@ import static org.springframework.cloud.loadbalancer.core.DiscoveryClientService * Tests for {@link DiscoveryClientServiceInstanceListSupplier}. * * @author Olga Maciaszek-Sharma + * @author Rod Catter */ class DiscoveryClientServiceInstanceListSupplierTests {