From 52c5b9e264740c90208885cb059e041af91c51ec Mon Sep 17 00:00:00 2001 From: Olga Maciaszek-Sharma Date: Mon, 16 Sep 2019 15:55:31 +0200 Subject: [PATCH] Switch ServiceInstanceListSupplier to use Flux. Add implementation with delayed list. Add subscriber-based resetting instances. --- ...scoveryClientServiceInstanceListSupplier.java | 14 ++++++++++---- .../core/PowerOfTwoChoicesPOCLoadBalancer.java | 16 +++++++++++----- .../core/RoundRobinListLoadBalancer.java | 1 + .../core/ServiceInstanceListSupplier.java | 4 ++-- 4 files changed, 24 insertions(+), 11 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 6d5100c7..7a3044ca 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 @@ -1,9 +1,10 @@ package org.springframework.cloud.loadbalancer.core; +import java.time.Duration; import java.util.List; import java.util.stream.Collectors; -import reactor.core.publisher.Mono; +import reactor.core.publisher.Flux; import org.springframework.cloud.client.discovery.DiscoveryClient; import org.springframework.core.env.Environment; @@ -26,14 +27,19 @@ public class DiscoveryClientServiceInstanceListSupplier implements ServiceInstan } @Override - public Mono> get() { - List instances = this.delegate + public Flux> get() { + //FIXME: sensible defaults + config + return Flux.just(getInstances()) + .delayElements(Duration.ofMinutes(5)); + } + + private List getInstances() { + return this.delegate .getInstances(this.serviceId) .stream() // switch to a more sensible conversion .map(serviceInstance -> (ConnectionTrackingServiceInstance) serviceInstance) .collect(Collectors.toList()); - return Mono.just(instances); } public String getServiceId() { diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/PowerOfTwoChoicesPOCLoadBalancer.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/PowerOfTwoChoicesPOCLoadBalancer.java index 738db6aa..442ef9da 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/PowerOfTwoChoicesPOCLoadBalancer.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/PowerOfTwoChoicesPOCLoadBalancer.java @@ -37,12 +37,18 @@ public class PowerOfTwoChoicesPOCLoadBalancer implements ReactorServiceInstanceL resetInstances(); } +// private void resetInstances() { +// Schedulers.fromExecutorService(Executors.newSingleThreadScheduledExecutor()) +// .schedulePeriodically(() -> instances = serviceInstanceListSupplier +// // maybe we don't have to block at all? +// // TODO: sensible interval defaults + config +// .getIfAvailable().get().next().block(), 0, 10, TimeUnit.MINUTES); +// } + private void resetInstances() { - Schedulers.fromExecutorService(Executors.newSingleThreadScheduledExecutor()) - .schedulePeriodically(() -> instances = serviceInstanceListSupplier - // maybe we don't have to block at all? - // TODO: sensible interval defaults + config - .getIfAvailable().get().block(), 0, 10, TimeUnit.MINUTES); + serviceInstanceListSupplier.getIfAvailable() + .get() + .subscribe(connectionTrackingServiceInstances -> instances = connectionTrackingServiceInstances); } // TODO: optimise diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/RoundRobinListLoadBalancer.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/RoundRobinListLoadBalancer.java index 6c845868..543b015b 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/RoundRobinListLoadBalancer.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/RoundRobinListLoadBalancer.java @@ -48,6 +48,7 @@ public class RoundRobinListLoadBalancer implements ReactorServiceInstanceLoadBal // TODO: move supplier to Request? ServiceInstanceListSupplier supplier = this.serviceInstanceListSupplier.getIfAvailable(); return supplier.get() + .next() .map(instances -> { if (instances.isEmpty()) { log.warn("No servers available for service: " + this.serviceId); diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplier.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplier.java index 8cd27330..1ad91211 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplier.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplier.java @@ -3,14 +3,14 @@ package org.springframework.cloud.loadbalancer.core; import java.util.List; import java.util.function.Supplier; -import reactor.core.publisher.Mono; +import reactor.core.publisher.Flux; import org.springframework.cloud.client.ServiceInstance; /** * @author Olga Maciaszek-Sharma */ -public interface ServiceInstanceListSupplier extends Supplier>> { +public interface ServiceInstanceListSupplier extends Supplier>> { String getServiceId(); }