From c0113c5235c079b8bae74bc4f56d71c89fcf1fe7 Mon Sep 17 00:00:00 2001 From: Tim Ysewyn Date: Thu, 3 Oct 2019 16:34:58 +0200 Subject: [PATCH] Support reactive discovery client for load balancing --- .../LoadBalancerClientConfiguration.java | 109 +++++++++++++----- ...veryClientServiceInstanceListSupplier.java | 19 ++- ...iscoveryClientServiceInstanceSupplier.java | 22 ++-- 3 files changed, 110 insertions(+), 40 deletions(-) diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/annotation/LoadBalancerClientConfiguration.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/annotation/LoadBalancerClientConfiguration.java index aaaa3cb3..1480ceee 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/annotation/LoadBalancerClientConfiguration.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/annotation/LoadBalancerClientConfiguration.java @@ -21,9 +21,12 @@ import org.springframework.boot.autoconfigure.condition.ConditionalOnBean; import org.springframework.boot.autoconfigure.condition.ConditionalOnMissingBean; import org.springframework.boot.context.properties.EnableConfigurationProperties; import org.springframework.cache.CacheManager; +import org.springframework.cloud.client.ConditionalOnBlockingDiscoveryEnabled; import org.springframework.cloud.client.ConditionalOnDiscoveryEnabled; +import org.springframework.cloud.client.ConditionalOnReactiveDiscoveryEnabled; import org.springframework.cloud.client.ServiceInstance; import org.springframework.cloud.client.discovery.DiscoveryClient; +import org.springframework.cloud.client.discovery.ReactiveDiscoveryClient; import org.springframework.cloud.loadbalancer.core.CachingServiceInstanceListSupplier; import org.springframework.cloud.loadbalancer.core.CachingServiceInstanceSupplier; import org.springframework.cloud.loadbalancer.core.DiscoveryClientServiceInstanceListSupplier; @@ -35,46 +38,20 @@ import org.springframework.cloud.loadbalancer.core.ServiceInstanceSupplier; import org.springframework.cloud.loadbalancer.support.LoadBalancerClientFactory; import org.springframework.context.annotation.Bean; import org.springframework.context.annotation.Configuration; +import org.springframework.core.annotation.Order; import org.springframework.core.env.Environment; /** * @author Spencer Gibb * @author Olga Maciaszek-Sharma + * @author Tim Ysewyn */ @Configuration @EnableConfigurationProperties @ConditionalOnDiscoveryEnabled public class LoadBalancerClientConfiguration { - @Bean - @ConditionalOnBean(DiscoveryClient.class) - @ConditionalOnMissingBean - public ServiceInstanceListSupplier discoveryClientServiceInstanceListSupplier( - DiscoveryClient discoveryClient, Environment env, - ObjectProvider cacheManager) { - DiscoveryClientServiceInstanceListSupplier delegate = new DiscoveryClientServiceInstanceListSupplier( - discoveryClient, env); - if (cacheManager.getIfAvailable() != null) { - return new CachingServiceInstanceListSupplier(delegate, - cacheManager.getIfAvailable()); - } - return delegate; - } - - @Bean - @ConditionalOnBean(DiscoveryClient.class) - @ConditionalOnMissingBean - public ServiceInstanceSupplier discoveryClientServiceInstanceSupplier( - DiscoveryClient discoveryClient, Environment env, - ObjectProvider cacheManager) { - DiscoveryClientServiceInstanceSupplier delegate = new DiscoveryClientServiceInstanceSupplier( - discoveryClient, env); - if (cacheManager.getIfAvailable() != null) { - return new CachingServiceInstanceSupplier(delegate, - cacheManager.getIfAvailable()); - } - return delegate; - } + private static final int REACTIVE_SERVICE_INSTANCE_SUPPLIER_ORDER = 193827465; @Bean @ConditionalOnMissingBean @@ -86,4 +63,78 @@ public class LoadBalancerClientConfiguration { ServiceInstanceListSupplier.class), name); } + @Configuration + @ConditionalOnReactiveDiscoveryEnabled + @Order(REACTIVE_SERVICE_INSTANCE_SUPPLIER_ORDER) + public static class ReactiveSupportConfiguration { + + @Bean + @ConditionalOnBean(ReactiveDiscoveryClient.class) + @ConditionalOnMissingBean + public ServiceInstanceListSupplier discoveryClientServiceInstanceListSupplier( + ReactiveDiscoveryClient discoveryClient, Environment env, + ObjectProvider cacheManager) { + DiscoveryClientServiceInstanceListSupplier delegate = new DiscoveryClientServiceInstanceListSupplier( + discoveryClient, env); + if (cacheManager.getIfAvailable() != null) { + return new CachingServiceInstanceListSupplier(delegate, + cacheManager.getIfAvailable()); + } + return delegate; + } + + @Bean + @ConditionalOnBean(ReactiveDiscoveryClient.class) + @ConditionalOnMissingBean + public ServiceInstanceSupplier discoveryClientServiceInstanceSupplier( + ReactiveDiscoveryClient discoveryClient, Environment env, + ObjectProvider cacheManager) { + DiscoveryClientServiceInstanceSupplier delegate = new DiscoveryClientServiceInstanceSupplier( + discoveryClient, env); + if (cacheManager.getIfAvailable() != null) { + return new CachingServiceInstanceSupplier(delegate, + cacheManager.getIfAvailable()); + } + return delegate; + } + + } + + @Configuration + @ConditionalOnBlockingDiscoveryEnabled + @Order(REACTIVE_SERVICE_INSTANCE_SUPPLIER_ORDER + 1) + public static class BlockingSupportConfiguration { + + @Bean + @ConditionalOnBean(DiscoveryClient.class) + @ConditionalOnMissingBean + public ServiceInstanceListSupplier discoveryClientServiceInstanceListSupplier( + DiscoveryClient discoveryClient, Environment env, + ObjectProvider cacheManager) { + DiscoveryClientServiceInstanceListSupplier delegate = new DiscoveryClientServiceInstanceListSupplier( + discoveryClient, env); + if (cacheManager.getIfAvailable() != null) { + return new CachingServiceInstanceListSupplier(delegate, + cacheManager.getIfAvailable()); + } + return delegate; + } + + @Bean + @ConditionalOnBean(DiscoveryClient.class) + @ConditionalOnMissingBean + public ServiceInstanceSupplier discoveryClientServiceInstanceSupplier( + DiscoveryClient discoveryClient, Environment env, + ObjectProvider cacheManager) { + DiscoveryClientServiceInstanceSupplier delegate = new DiscoveryClientServiceInstanceSupplier( + discoveryClient, env); + if (cacheManager.getIfAvailable() != null) { + return new CachingServiceInstanceSupplier(delegate, + cacheManager.getIfAvailable()); + } + return delegate; + } + + } + } 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 2ba4161f..eb4a7fa1 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 @@ -19,9 +19,11 @@ package org.springframework.cloud.loadbalancer.core; import java.util.List; import reactor.core.publisher.Flux; +import reactor.core.scheduler.Schedulers; import org.springframework.cloud.client.ServiceInstance; import org.springframework.cloud.client.discovery.DiscoveryClient; +import org.springframework.cloud.client.discovery.ReactiveDiscoveryClient; import org.springframework.core.env.Environment; import static org.springframework.cloud.loadbalancer.support.LoadBalancerClientFactory.PROPERTY_NAME; @@ -31,19 +33,28 @@ import static org.springframework.cloud.loadbalancer.support.LoadBalancerClientF * * @author Spencer Gibb * @author Olga Maciaszek-Sharma + * @author Tim Ysewyn * @since 2.2.0 */ public class DiscoveryClientServiceInstanceListSupplier implements ServiceInstanceListSupplier { - private final DiscoveryClient delegate; - private final String serviceId; + private final Flux serviceInstances; + public DiscoveryClientServiceInstanceListSupplier(DiscoveryClient delegate, Environment environment) { - this.delegate = delegate; this.serviceId = environment.getProperty(PROPERTY_NAME); + this.serviceInstances = Flux + .defer(() -> Flux.fromIterable(delegate.getInstances(serviceId))) + .subscribeOn(Schedulers.boundedElastic()); + } + + public DiscoveryClientServiceInstanceListSupplier(ReactiveDiscoveryClient delegate, + Environment environment) { + this.serviceId = environment.getProperty(PROPERTY_NAME); + this.serviceInstances = delegate.getInstances(serviceId); } @Override @@ -53,7 +64,7 @@ public class DiscoveryClientServiceInstanceListSupplier @Override public Flux> get() { - return Flux.defer(() -> Flux.just(delegate.getInstances(serviceId))); + return serviceInstances.collectList().flux(); } } diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DiscoveryClientServiceInstanceSupplier.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DiscoveryClientServiceInstanceSupplier.java index b6a8b53a..fbf96a10 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DiscoveryClientServiceInstanceSupplier.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DiscoveryClientServiceInstanceSupplier.java @@ -16,12 +16,12 @@ package org.springframework.cloud.loadbalancer.core; -import java.util.List; - import reactor.core.publisher.Flux; +import reactor.core.scheduler.Schedulers; import org.springframework.cloud.client.ServiceInstance; import org.springframework.cloud.client.discovery.DiscoveryClient; +import org.springframework.cloud.client.discovery.ReactiveDiscoveryClient; import org.springframework.core.env.Environment; import static org.springframework.cloud.loadbalancer.support.LoadBalancerClientFactory.PROPERTY_NAME; @@ -29,24 +29,32 @@ import static org.springframework.cloud.loadbalancer.support.LoadBalancerClientF /** * @deprecated Use {@link DiscoveryClientServiceInstanceListSupplier} instead. * @author Spencer Gibb + * @author Tim Ysewyn */ @Deprecated public class DiscoveryClientServiceInstanceSupplier implements ServiceInstanceSupplier { - private final DiscoveryClient delegate; - private final String serviceId; + private final Flux serviceInstances; + public DiscoveryClientServiceInstanceSupplier(DiscoveryClient delegate, Environment environment) { - this.delegate = delegate; this.serviceId = environment.getProperty(PROPERTY_NAME); + this.serviceInstances = Flux + .defer(() -> Flux.fromIterable(delegate.getInstances(serviceId))) + .subscribeOn(Schedulers.boundedElastic()); + } + + public DiscoveryClientServiceInstanceSupplier(ReactiveDiscoveryClient delegate, + Environment environment) { + this.serviceId = environment.getProperty(PROPERTY_NAME); + this.serviceInstances = delegate.getInstances(serviceId); } @Override public Flux get() { - List instances = this.delegate.getInstances(this.serviceId); - return Flux.fromIterable(instances); + return this.serviceInstances; } public String getServiceId() {