From 3c310a819dd2cf326fc6f416813f6a133c143459 Mon Sep 17 00:00:00 2001 From: Olga Maciaszek-Sharma Date: Fri, 13 Sep 2019 15:15:29 +0200 Subject: [PATCH] Draft POC. --- .../ConnectionTrackingServiceInstance.java | 14 ++++ ...veryClientServiceInstanceListSupplier.java | 42 ++++++++++++ .../PowerOfTwoChoicesPOCLoadBalancer.java | 67 +++++++++++++++++++ .../core/RoundRobinListLoadBalancer.java | 64 ++++++++++++++++++ .../core/ServiceInstanceListSupplier.java | 16 +++++ 5 files changed, 203 insertions(+) create mode 100644 spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ConnectionTrackingServiceInstance.java create mode 100644 spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DiscoveryClientServiceInstanceListSupplier.java create mode 100644 spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/PowerOfTwoChoicesPOCLoadBalancer.java create mode 100644 spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/RoundRobinListLoadBalancer.java create mode 100644 spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplier.java diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ConnectionTrackingServiceInstance.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ConnectionTrackingServiceInstance.java new file mode 100644 index 00000000..a850a9a6 --- /dev/null +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ConnectionTrackingServiceInstance.java @@ -0,0 +1,14 @@ +package org.springframework.cloud.loadbalancer.core; + +import org.springframework.cloud.client.ServiceInstance; + +/** + * @author Olga Maciaszek-Sharma + */ +public interface ConnectionTrackingServiceInstance extends ServiceInstance { + + int getConnectionCount(); + + void addConnection(); + +} 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 new file mode 100644 index 00000000..6d5100c7 --- /dev/null +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DiscoveryClientServiceInstanceListSupplier.java @@ -0,0 +1,42 @@ +package org.springframework.cloud.loadbalancer.core; + +import java.util.List; +import java.util.stream.Collectors; + +import reactor.core.publisher.Mono; + +import org.springframework.cloud.client.discovery.DiscoveryClient; +import org.springframework.core.env.Environment; + +import static org.springframework.cloud.loadbalancer.support.LoadBalancerClientFactory.PROPERTY_NAME; + +/** + * @author Olga Maciaszek-Sharma + */ +public class DiscoveryClientServiceInstanceListSupplier implements ServiceInstanceListSupplier { + + private final DiscoveryClient delegate; + + private final String serviceId; + + public DiscoveryClientServiceInstanceListSupplier(DiscoveryClient delegate, + Environment environment) { + this.delegate = delegate; + this.serviceId = environment.getProperty(PROPERTY_NAME); + } + + @Override + public Mono> get() { + List instances = 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() { + return this.serviceId; + } +} 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 new file mode 100644 index 00000000..738db6aa --- /dev/null +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/PowerOfTwoChoicesPOCLoadBalancer.java @@ -0,0 +1,67 @@ +package org.springframework.cloud.loadbalancer.core; + +import java.util.Collections; +import java.util.List; +import java.util.concurrent.Executors; +import java.util.concurrent.TimeUnit; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import reactor.core.publisher.Mono; +import reactor.core.scheduler.Schedulers; + +import org.springframework.beans.factory.ObjectProvider; +import org.springframework.cloud.client.ServiceInstance; +import org.springframework.cloud.client.loadbalancer.reactive.DefaultResponse; +import org.springframework.cloud.client.loadbalancer.reactive.EmptyResponse; +import org.springframework.cloud.client.loadbalancer.reactive.Request; +import org.springframework.cloud.client.loadbalancer.reactive.Response; + +/** + * @author Olga Maciaszek-Sharma + */ +public class PowerOfTwoChoicesPOCLoadBalancer implements ReactorServiceInstanceLoadBalancer { + + private static final Log log = LogFactory.getLog(RoundRobinLoadBalancer.class); + + private final ObjectProvider> serviceInstanceListSupplier; + + private final String serviceId; + + private List instances; + + public PowerOfTwoChoicesPOCLoadBalancer(String serviceId, + ObjectProvider> serviceInstanceListSupplier) { + this.serviceId = serviceId; + this.serviceInstanceListSupplier = serviceInstanceListSupplier; + 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().block(), 0, 10, TimeUnit.MINUTES); + } + + // TODO: optimise + @Override + public Mono> choose(Request request) { + if (instances.isEmpty()) { + log.warn("No servers available for service: " + this.serviceId); + return Mono.just(new EmptyResponse()); + } + if (instances.size() == 1 || instances.get(0).getConnectionCount() < instances + .get(1) + .getConnectionCount()) { + instances.get(0).addConnection(); + return Mono.just(new DefaultResponse(instances.get(0))); + } + Collections.shuffle(instances); + + instances.get(1).addConnection(); + return Mono.just(new DefaultResponse(instances.get(1))); + } + +} 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 new file mode 100644 index 00000000..6c845868 --- /dev/null +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/RoundRobinListLoadBalancer.java @@ -0,0 +1,64 @@ +package org.springframework.cloud.loadbalancer.core; + +import java.util.Random; +import java.util.concurrent.atomic.AtomicInteger; + +import org.apache.commons.logging.Log; +import org.apache.commons.logging.LogFactory; +import reactor.core.publisher.Mono; + +import org.springframework.beans.factory.ObjectProvider; +import org.springframework.cloud.client.ServiceInstance; +import org.springframework.cloud.client.loadbalancer.reactive.DefaultResponse; +import org.springframework.cloud.client.loadbalancer.reactive.EmptyResponse; +import org.springframework.cloud.client.loadbalancer.reactive.Request; +import org.springframework.cloud.client.loadbalancer.reactive.Response; + +/** + * @author Olga Maciaszek-Sharma + */ +public class RoundRobinListLoadBalancer implements ReactorServiceInstanceLoadBalancer { + + private static final Log log = LogFactory.getLog(RoundRobinLoadBalancer.class); + + private final AtomicInteger position; + + private final ObjectProvider> serviceInstanceListSupplier; + + private final String serviceId; + + public RoundRobinListLoadBalancer(String serviceId, + ObjectProvider> serviceInstanceListSupplier) { + this(serviceId, serviceInstanceListSupplier, new Random().nextInt(1000)); + } + + public RoundRobinListLoadBalancer(String serviceId, + ObjectProvider> serviceInstanceListSupplier, + int seedPosition) { + this.serviceId = serviceId; + this.serviceInstanceListSupplier = serviceInstanceListSupplier; + this.position = new AtomicInteger(seedPosition); + } + + @Override + // see original + // https://github.com/Netflix/ocelli/blob/master/ocelli-core/ + // src/main/java/netflix/ocelli/loadbalancer/RoundRobinLoadBalancer.java + public Mono> choose(Request request) { + // TODO: move supplier to Request? + ServiceInstanceListSupplier supplier = this.serviceInstanceListSupplier.getIfAvailable(); + return supplier.get() + .map(instances -> { + if (instances.isEmpty()) { + log.warn("No servers available for service: " + this.serviceId); + return new EmptyResponse(); + } + // TODO: enforce order? + int pos = Math.abs(this.position.incrementAndGet()); + + ServiceInstance instance = instances.get(pos % instances.size()); + + return new DefaultResponse(instance); + }); + } +} 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 new file mode 100644 index 00000000..8cd27330 --- /dev/null +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplier.java @@ -0,0 +1,16 @@ +package org.springframework.cloud.loadbalancer.core; + +import java.util.List; +import java.util.function.Supplier; + +import reactor.core.publisher.Mono; + +import org.springframework.cloud.client.ServiceInstance; + +/** + * @author Olga Maciaszek-Sharma + */ +public interface ServiceInstanceListSupplier extends Supplier>> { + + String getServiceId(); +}