From 354a894606217fd466a30106c5ea736abcdf44ab Mon Sep 17 00:00:00 2001 From: Olga Maciaszek-Sharma Date: Wed, 28 Jun 2023 12:28:44 +0200 Subject: [PATCH 1/3] Call get request on delegates (#1250) --- .../main/asciidoc/spring-cloud-commons.adoc | 26 +++++++++-- .../loadbalancer/LoadBalancerProperties.java | 43 +++++++++++++++++++ .../LoadBalancerClientConfiguration.java | 22 +++++----- ...DelegatingServiceInstanceListSupplier.java | 4 +- ...PreferenceServiceInstanceListSupplier.java | 19 ++++++++ .../ServiceInstanceListSupplierBuilder.java | 32 ++++++-------- ...PreferenceServiceInstanceListSupplier.java | 21 +++++++++ .../LoadBalancerClientConfigurationTests.java | 8 ++-- ...renceServiceInstanceListSupplierTests.java | 33 ++++++++++++-- ...rviceInstanceListSupplierBuilderTests.java | 6 +-- ...renceServiceInstanceListSupplierTests.java | 33 +++++++++++++- 11 files changed, 200 insertions(+), 47 deletions(-) diff --git a/docs/src/main/asciidoc/spring-cloud-commons.adoc b/docs/src/main/asciidoc/spring-cloud-commons.adoc index 97fc2423..ca20aa33 100644 --- a/docs/src/main/asciidoc/spring-cloud-commons.adoc +++ b/docs/src/main/asciidoc/spring-cloud-commons.adoc @@ -933,6 +933,11 @@ to `false`. WARNING: Although the basic, non-cached, implementation is useful for prototyping and testing, it's much less efficient than the cached versions, so we recommend always using the cached version in production. If the caching is already done by the `DiscoveryClient` implementation, for example `EurekaDiscoveryClient`, the load-balancer caching should be disabled to prevent double caching. +==== + +NOTE: When you create your own configuration, if you use `CachingServiceInstanceListSupplier` make sure to place it in the hierarchy directly after the supplier that retrieves the instances over the network, for example, `DiscoveryClientServiceInstanceListSupplier`, before any other filtering suppliers. + +==== === Zone-Based Load-Balancing To enable zone-based load-balancing, we provide the `ZonePreferenceServiceInstanceListSupplier`. @@ -950,7 +955,7 @@ If the zone is `null` or there are no instances within the same zone, it returns In order to use the zone-based load-balancing approach, you will have to instantiate a `ZonePreferenceServiceInstanceListSupplier` bean in a <>. We use delegates to work with `ServiceInstanceListSupplier` beans. -We suggest passing a `DiscoveryClientServiceInstanceListSupplier` delegate in the constructor of `ZonePreferenceServiceInstanceListSupplier` and, in turn, wrapping the latter with a `CachingServiceInstanceListSupplier` to leverage <>. +We suggest using a `DiscoveryClientServiceInstanceListSupplier` delegate, wrapping it with a `CachingServiceInstanceListSupplier` to leverage <>, and then passing the resulting bean in the constructor of `ZonePreferenceServiceInstanceListSupplier`. You can use this sample configuration to set it up: @@ -964,8 +969,8 @@ public class CustomLoadBalancerConfiguration { ConfigurableApplicationContext context) { return ServiceInstanceListSupplier.builder() .withDiscoveryClient() + .withCaching() .withZonePreference() - .withCaching() .build(context); } } @@ -1026,6 +1031,12 @@ You can also pass your own `WebClient` or `RestTemplate` instance to be used for WARNING: `HealthCheckServiceInstanceListSupplier` has its own caching mechanism based on Reactor Flux `replay()`. Therefore, if it's being used, you may want to skip wrapping that supplier with `CachingServiceInstanceListSupplier`. +==== + +NOTE: When you create your own configuration, `HealthCheckServiceInstanceListSupplier`, make sure to place it in the hierarchy directly after the supplier that retrieves the instances over the network, for example, `DiscoveryClientServiceInstanceListSupplier`, before any other filtering suppliers. + +==== + === Same instance preference for LoadBalancer You can set up the LoadBalancer in such a way that it prefers the instance that was previously selected, if that instance is available. @@ -1110,8 +1121,8 @@ public class CustomLoadBalancerConfiguration { ConfigurableApplicationContext context) { return ServiceInstanceListSupplier.builder() .withDiscoveryClient() + .withCaching() .withHints() - .withCaching() .build(context); } } @@ -1221,11 +1232,18 @@ public class MyConfiguration { } } ---- +==== NOTE: The classes you pass as `@LoadBalancerClient` or `@LoadBalancerClients` configuration arguments should either not be annotated with `@Configuration` or be outside component scan scope. ==== +==== + +NOTE: When you create your own configuration, if you use `CachingServiceInstanceListSupplier` or `HealthCheckServiceInstanceListSupplier`, makes sure to use one of them, not both, and make sure to place it in the hierarchy directly after the supplier that retrieves the instances over the network, for example, `DiscoveryClientServiceInstanceListSupplier`, before any other filtering suppliers. + +==== + [[loadbalancer-lifecycle]] === Spring Cloud LoadBalancer Lifecycle @@ -1301,6 +1319,8 @@ The per-client configuration properties work for most of the properties, apart f NOTE: For the properties where maps where already used, where you can specify a different value per-client without using the `clients` keyword (for example, `hints`, `health-check.path`), we have kept that behaviour in order to keep the library backwards compatible. It will be modified in the next major release. +NOTE: Starting with `3.1.7` in `2021.0.x` release train, `4.0.4` in `2022.0.x` release train and `4.1.0` in the `2023.0.x` release train, we have introduced the `callGetWithRequestOnDelegates` flag in `LoadBalancerProperties`. If this flag is set to `true`, `ServiceInstanceListSupplier#get(Request request)` method will be implemented to call `delegate.get(request)` in classes assignable from `DelegatingServiceInstanceListSupplier` that don't already implement that method, with the exclusion of `CachingServiceInstanceListSupplier` and `HealthCheckServiceInstanceListSupplier`, which should be placed in the instance supplier hierarchy directly after the supplier performing instance retrieval over the network, before any request-based filtering is done. For `3.1.x` and `4.0.x` the flag is set to `false` by default, and since `4.1.0` it's going to be set to `true` by default. + == Spring Cloud Circuit Breaker include::spring-cloud-circuitbreaker.adoc[leveloffset=+1] diff --git a/spring-cloud-commons/src/main/java/org/springframework/cloud/client/loadbalancer/LoadBalancerProperties.java b/spring-cloud-commons/src/main/java/org/springframework/cloud/client/loadbalancer/LoadBalancerProperties.java index 4b719cc5..5537dde1 100644 --- a/spring-cloud-commons/src/main/java/org/springframework/cloud/client/loadbalancer/LoadBalancerProperties.java +++ b/spring-cloud-commons/src/main/java/org/springframework/cloud/client/loadbalancer/LoadBalancerProperties.java @@ -73,6 +73,19 @@ public class LoadBalancerProperties { */ private boolean useRawStatusCodeInResponseData; + /** + * If this flag is set to {@code true}, + * {@code ServiceInstanceListSupplier#get(Request request)} method will be implemented + * to call {@code delegate.get(request)} in classes assignable from + * {@code DelegatingServiceInstanceListSupplier} that don't already implement that + * method, with the exclusion of {@code CachingServiceInstanceListSupplier} and + * {@code HealthCheckServiceInstanceListSupplier}, which should be placed in the + * instance supplier hierarchy directly after the supplier performing instance + * retrieval over the network, before any request-based filtering is done. Note: in + * 4.1, this behaviour will become the default + */ + private boolean callGetWithRequestOnDelegates; + public HealthCheck getHealthCheck() { return healthCheck; } @@ -134,6 +147,36 @@ public class LoadBalancerProperties { this.useRawStatusCodeInResponseData = useRawStatusCodeInResponseData; } + /** + * If this flag is set to {@code true}, + * {@code ServiceInstanceListSupplier#get(Request request)} method will be implemented + * to call {@code delegate.get(request)} in classes assignable from + * {@code DelegatingServiceInstanceListSupplier} that don't already implement that + * method, with the exclusion of {@code CachingServiceInstanceListSupplier} and + * {@code HealthCheckServiceInstanceListSupplier}, which should be placed in the + * instance supplier hierarchy directly after the supplier performing instance + * retrieval over the network, before any request-based filtering is done. Note: in + * 4.1, this behaviour will become the default + */ + public boolean isCallGetWithRequestOnDelegates() { + return callGetWithRequestOnDelegates; + } + + /** + * If this flag is set to {@code true}, + * {@code ServiceInstanceListSupplier#get(Request request)} method will be implemented + * to call {@code delegate.get(request)} in classes assignable from + * {@code DelegatingServiceInstanceListSupplier} that don't already implement that + * method, with the exclusion of {@code CachingServiceInstanceListSupplier} and + * {@code HealthCheckServiceInstanceListSupplier}, which should be placed in the + * instance supplier hierarchy directly after the supplier performing instance + * retrieval over the network, before any request-based filtering is done. Note: in + * 4.1, this behaviour will become the default + */ + public void setCallGetWithRequestOnDelegates(boolean callGetWithRequestOnDelegates) { + this.callGetWithRequestOnDelegates = callGetWithRequestOnDelegates; + } + public static class StickySession { /** 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 7c2cc898..4f203f07 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 @@ -93,7 +93,7 @@ public class LoadBalancerClientConfiguration { @Conditional(ZonePreferenceConfigurationCondition.class) public ServiceInstanceListSupplier zonePreferenceDiscoveryClientServiceInstanceListSupplier( ConfigurableApplicationContext context) { - return ServiceInstanceListSupplier.builder().withDiscoveryClient().withZonePreference().withCaching() + return ServiceInstanceListSupplier.builder().withDiscoveryClient().withCaching().withZonePreference() .build(context); } @@ -119,8 +119,8 @@ public class LoadBalancerClientConfiguration { @Conditional(RequestBasedStickySessionConfigurationCondition.class) public ServiceInstanceListSupplier requestBasedStickySessionDiscoveryClientServiceInstanceListSupplier( ConfigurableApplicationContext context) { - return ServiceInstanceListSupplier.builder().withDiscoveryClient().withRequestBasedStickySession() - .withCaching().build(context); + return ServiceInstanceListSupplier.builder().withDiscoveryClient().withCaching() + .withRequestBasedStickySession().build(context); } @Bean @@ -129,8 +129,8 @@ public class LoadBalancerClientConfiguration { @Conditional(SameInstancePreferenceConfigurationCondition.class) public ServiceInstanceListSupplier sameInstancePreferenceServiceInstanceListSupplier( ConfigurableApplicationContext context) { - return ServiceInstanceListSupplier.builder().withDiscoveryClient().withSameInstancePreference() - .withCaching().build(context); + return ServiceInstanceListSupplier.builder().withDiscoveryClient().withCaching() + .withSameInstancePreference().build(context); } } @@ -155,8 +155,8 @@ public class LoadBalancerClientConfiguration { @Conditional(ZonePreferenceConfigurationCondition.class) public ServiceInstanceListSupplier zonePreferenceDiscoveryClientServiceInstanceListSupplier( ConfigurableApplicationContext context) { - return ServiceInstanceListSupplier.builder().withBlockingDiscoveryClient().withZonePreference() - .withCaching().build(context); + return ServiceInstanceListSupplier.builder().withBlockingDiscoveryClient().withCaching() + .withZonePreference().build(context); } @Bean @@ -175,8 +175,8 @@ public class LoadBalancerClientConfiguration { @Conditional(RequestBasedStickySessionConfigurationCondition.class) public ServiceInstanceListSupplier requestBasedStickySessionDiscoveryClientServiceInstanceListSupplier( ConfigurableApplicationContext context) { - return ServiceInstanceListSupplier.builder().withBlockingDiscoveryClient().withRequestBasedStickySession() - .withCaching().build(context); + return ServiceInstanceListSupplier.builder().withBlockingDiscoveryClient().withCaching() + .withRequestBasedStickySession().build(context); } @Bean @@ -185,8 +185,8 @@ public class LoadBalancerClientConfiguration { @Conditional(SameInstancePreferenceConfigurationCondition.class) public ServiceInstanceListSupplier sameInstancePreferenceServiceInstanceListSupplier( ConfigurableApplicationContext context) { - return ServiceInstanceListSupplier.builder().withBlockingDiscoveryClient().withSameInstancePreference() - .withCaching().build(context); + return ServiceInstanceListSupplier.builder().withBlockingDiscoveryClient().withCaching() + .withSameInstancePreference().build(context); } } diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DelegatingServiceInstanceListSupplier.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DelegatingServiceInstanceListSupplier.java index 50ba6a18..af791aa4 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DelegatingServiceInstanceListSupplier.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DelegatingServiceInstanceListSupplier.java @@ -38,12 +38,12 @@ public abstract class DelegatingServiceInstanceListSupplier } public ServiceInstanceListSupplier getDelegate() { - return this.delegate; + return delegate; } @Override public String getServiceId() { - return this.delegate.getServiceId(); + return delegate.getServiceId(); } @Override diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplier.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplier.java index d198051a..8317d9d9 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplier.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplier.java @@ -24,6 +24,8 @@ import org.apache.commons.logging.LogFactory; import reactor.core.publisher.Flux; import org.springframework.cloud.client.ServiceInstance; +import org.springframework.cloud.client.loadbalancer.Request; +import org.springframework.cloud.client.loadbalancer.reactive.ReactiveLoadBalancer; /** * An implementation of {@link ServiceInstanceListSupplier} that selects the previously @@ -39,10 +41,19 @@ public class SameInstancePreferenceServiceInstanceListSupplier extends Delegatin private ServiceInstance previouslyReturnedInstance; + private boolean callGetWithRequestOnDelegates; + public SameInstancePreferenceServiceInstanceListSupplier(ServiceInstanceListSupplier delegate) { super(delegate); } + public SameInstancePreferenceServiceInstanceListSupplier(ServiceInstanceListSupplier delegate, + ReactiveLoadBalancer.Factory loadBalancerClientFactory) { + super(delegate); + callGetWithRequestOnDelegates = loadBalancerClientFactory.getProperties(getServiceId()) + .isCallGetWithRequestOnDelegates(); + } + @Override public String getServiceId() { return delegate.getServiceId(); @@ -53,6 +64,14 @@ public class SameInstancePreferenceServiceInstanceListSupplier extends Delegatin return delegate.get().map(this::filteredBySameInstancePreference); } + @Override + public Flux> get(Request request) { + if (callGetWithRequestOnDelegates) { + return delegate.get(request).map(this::filteredBySameInstancePreference); + } + return get(); + } + private List filteredBySameInstancePreference(List serviceInstances) { if (previouslyReturnedInstance != null && serviceInstances.contains(previouslyReturnedInstance)) { if (LOG.isDebugEnabled()) { diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilder.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilder.java index 6a19838e..0168d98b 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilder.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilder.java @@ -57,8 +57,6 @@ public final class ServiceInstanceListSupplierBuilder { private Creator baseCreator; - private DelegateCreator cachingCreator; - private final List creators = new ArrayList<>(); ServiceInstanceListSupplierBuilder() { @@ -148,8 +146,10 @@ public final class ServiceInstanceListSupplierBuilder { * @return the {@link ServiceInstanceListSupplierBuilder} object */ public ServiceInstanceListSupplierBuilder withSameInstancePreference() { - DelegateCreator creator = (context, - delegate) -> new SameInstancePreferenceServiceInstanceListSupplier(delegate); + DelegateCreator creator = (context, delegate) -> { + LoadBalancerClientFactory loadBalancerClientFactory = context.getBean(LoadBalancerClientFactory.class); + return new SameInstancePreferenceServiceInstanceListSupplier(delegate, loadBalancerClientFactory); + }; this.creators.add(creator); return this; } @@ -191,8 +191,9 @@ public final class ServiceInstanceListSupplierBuilder { */ public ServiceInstanceListSupplierBuilder withZonePreference() { DelegateCreator creator = (context, delegate) -> { + LoadBalancerClientFactory loadBalancerClientFactory = context.getBean(LoadBalancerClientFactory.class); LoadBalancerZoneConfig zoneConfig = context.getBean(LoadBalancerZoneConfig.class); - return new ZonePreferenceServiceInstanceListSupplier(delegate, zoneConfig); + return new ZonePreferenceServiceInstanceListSupplier(delegate, zoneConfig, loadBalancerClientFactory); }; this.creators.add(creator); return this; @@ -206,8 +207,9 @@ public final class ServiceInstanceListSupplierBuilder { */ public ServiceInstanceListSupplierBuilder withZonePreference(String zoneName) { DelegateCreator creator = (context, delegate) -> { + LoadBalancerClientFactory loadBalancerClientFactory = context.getBean(LoadBalancerClientFactory.class); LoadBalancerZoneConfig zoneConfig = new LoadBalancerZoneConfig(zoneName); - return new ZonePreferenceServiceInstanceListSupplier(delegate, zoneConfig); + return new ZonePreferenceServiceInstanceListSupplier(delegate, zoneConfig, loadBalancerClientFactory); }; this.creators.add(creator); return this; @@ -228,19 +230,15 @@ public final class ServiceInstanceListSupplierBuilder { } /** - * If {@link LoadBalancerCacheManager} is available in the context, wraps created - * {@link ServiceInstanceListSupplier} hierarchy with a - * {@link CachingServiceInstanceListSupplier} instance to provide a caching mechanism - * for service instances. Uses {@link ObjectProvider} to lazily resolve + * If {@link LoadBalancerCacheManager} is available in the context, adds a + * {@link CachingServiceInstanceListSupplier} instance to the + * {@link ServiceInstanceListSupplier} hierarchy to provide a caching mechanism for + * service instances. Uses {@link ObjectProvider} to lazily resolve * {@link LoadBalancerCacheManager}. * @return the {@link ServiceInstanceListSupplierBuilder} object */ public ServiceInstanceListSupplierBuilder withCaching() { - if (cachingCreator != null && LOG.isWarnEnabled()) { - LOG.warn( - "Overriding a previously set cachingCreator with a CachingServiceInstanceListSupplier-based cachingCreator."); - } - this.cachingCreator = (context, delegate) -> { + DelegateCreator creator = (context, delegate) -> { ObjectProvider cacheManagerProvider = context .getBeanProvider(LoadBalancerCacheManager.class); if (cacheManagerProvider.getIfAvailable() != null) { @@ -251,6 +249,7 @@ public final class ServiceInstanceListSupplierBuilder { } return delegate; }; + creators.add(creator); return this; } @@ -297,9 +296,6 @@ public final class ServiceInstanceListSupplierBuilder { supplier = creator.apply(context, supplier); } - if (this.cachingCreator != null) { - supplier = this.cachingCreator.apply(context, supplier); - } return supplier; } diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ZonePreferenceServiceInstanceListSupplier.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ZonePreferenceServiceInstanceListSupplier.java index 83ad66a7..26f1c7da 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ZonePreferenceServiceInstanceListSupplier.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ZonePreferenceServiceInstanceListSupplier.java @@ -23,6 +23,8 @@ import java.util.Map; import reactor.core.publisher.Flux; import org.springframework.cloud.client.ServiceInstance; +import org.springframework.cloud.client.loadbalancer.Request; +import org.springframework.cloud.client.loadbalancer.reactive.ReactiveLoadBalancer; import org.springframework.cloud.loadbalancer.config.LoadBalancerZoneConfig; /** @@ -43,17 +45,36 @@ public class ZonePreferenceServiceInstanceListSupplier extends DelegatingService private String zone; + private boolean callGetWithRequestOnDelegates; + public ZonePreferenceServiceInstanceListSupplier(ServiceInstanceListSupplier delegate, LoadBalancerZoneConfig zoneConfig) { super(delegate); this.zoneConfig = zoneConfig; } + public ZonePreferenceServiceInstanceListSupplier(ServiceInstanceListSupplier delegate, + LoadBalancerZoneConfig zoneConfig, + ReactiveLoadBalancer.Factory loadBalancerClientFactory) { + super(delegate); + this.zoneConfig = zoneConfig; + callGetWithRequestOnDelegates = loadBalancerClientFactory.getProperties(getServiceId()) + .isCallGetWithRequestOnDelegates(); + } + @Override public Flux> get() { return getDelegate().get().map(this::filteredByZone); } + @Override + public Flux> get(Request request) { + if (callGetWithRequestOnDelegates) { + return getDelegate().get(request).map(this::filteredByZone); + } + return get(); + } + private List filteredByZone(List serviceInstances) { if (zone == null) { zone = zoneConfig.getZone(); diff --git a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/annotation/LoadBalancerClientConfigurationTests.java b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/annotation/LoadBalancerClientConfigurationTests.java index 9ac93283..ebd12fea 100644 --- a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/annotation/LoadBalancerClientConfigurationTests.java +++ b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/annotation/LoadBalancerClientConfigurationTests.java @@ -91,10 +91,10 @@ class LoadBalancerClientConfigurationTests { reactiveDiscoveryClientRunner.withPropertyValues("spring.cloud.loadbalancer.configurations=zone-preference") .run(context -> { ServiceInstanceListSupplier supplier = context.getBean(ServiceInstanceListSupplier.class); - then(supplier).isInstanceOf(CachingServiceInstanceListSupplier.class); + then(supplier).isInstanceOf(ZonePreferenceServiceInstanceListSupplier.class); ServiceInstanceListSupplier delegate = ((DelegatingServiceInstanceListSupplier) supplier) .getDelegate(); - then(delegate).isInstanceOf(ZonePreferenceServiceInstanceListSupplier.class); + then(delegate).isInstanceOf(CachingServiceInstanceListSupplier.class); ServiceInstanceListSupplier secondDelegate = ((DelegatingServiceInstanceListSupplier) delegate) .getDelegate(); then(secondDelegate).isInstanceOf(DiscoveryClientServiceInstanceListSupplier.class); @@ -119,10 +119,10 @@ class LoadBalancerClientConfigurationTests { .withPropertyValues("spring.cloud.loadbalancer.configurations=request-based-sticky-session") .run(context -> { ServiceInstanceListSupplier supplier = context.getBean(ServiceInstanceListSupplier.class); - then(supplier).isInstanceOf(CachingServiceInstanceListSupplier.class); + then(supplier).isInstanceOf(RequestBasedStickySessionServiceInstanceListSupplier.class); ServiceInstanceListSupplier delegate = ((DelegatingServiceInstanceListSupplier) supplier) .getDelegate(); - then(delegate).isInstanceOf(RequestBasedStickySessionServiceInstanceListSupplier.class); + then(delegate).isInstanceOf(CachingServiceInstanceListSupplier.class); ServiceInstanceListSupplier secondDelegate = ((DelegatingServiceInstanceListSupplier) delegate) .getDelegate(); then(secondDelegate).isInstanceOf(DiscoveryClientServiceInstanceListSupplier.class); diff --git a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplierTests.java b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplierTests.java index fd66a251..bf6e260b 100644 --- a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplierTests.java +++ b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplierTests.java @@ -19,13 +19,20 @@ package org.springframework.cloud.loadbalancer.core; import java.util.Arrays; import java.util.List; +import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import reactor.core.publisher.Flux; import org.springframework.cloud.client.DefaultServiceInstance; import org.springframework.cloud.client.ServiceInstance; +import org.springframework.cloud.client.loadbalancer.DefaultRequest; +import org.springframework.cloud.client.loadbalancer.DefaultRequestContext; +import org.springframework.cloud.client.loadbalancer.LoadBalancerProperties; +import org.springframework.cloud.client.loadbalancer.Request; +import org.springframework.cloud.loadbalancer.support.LoadBalancerClientFactory; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; @@ -39,8 +46,9 @@ class SameInstancePreferenceServiceInstanceListSupplierTests { private final DiscoveryClientServiceInstanceListSupplier delegate = mock( DiscoveryClientServiceInstanceListSupplier.class); - private final SameInstancePreferenceServiceInstanceListSupplier supplier = new SameInstancePreferenceServiceInstanceListSupplier( - delegate); + private final LoadBalancerClientFactory loadBalancerClientFactory = mock(LoadBalancerClientFactory.class); + + private SameInstancePreferenceServiceInstanceListSupplier supplier; private final ServiceInstance first = serviceInstance("test-1"); @@ -48,6 +56,14 @@ class SameInstancePreferenceServiceInstanceListSupplierTests { private final ServiceInstance third = serviceInstance("test-3"); + @BeforeEach + void setUp() { + LoadBalancerProperties properties = new LoadBalancerProperties(); + properties.setCallGetWithRequestOnDelegates(true); + when(loadBalancerClientFactory.getProperties(any())).thenReturn(properties); + supplier = new SameInstancePreferenceServiceInstanceListSupplier(delegate, loadBalancerClientFactory); + } + @Test void shouldReturnPreviouslySelectedInstanceIfAvailable() { when(delegate.get()).thenReturn(Flux.just(Arrays.asList(first, second, third))); @@ -69,7 +85,7 @@ class SameInstancePreferenceServiceInstanceListSupplierTests { } @Test - void shouldReturnAllInstancesFromDelegateIfPreviouslySelectedInstanceIfAvailable() { + void shouldReturnAllInstancesFromDelegateIfPreviouslySelectedInstanceIsNotAvailable() { when(delegate.get()).thenReturn(Flux.just(Arrays.asList(second, third))); supplier.selectedServiceInstance(first); @@ -78,6 +94,17 @@ class SameInstancePreferenceServiceInstanceListSupplierTests { assertThat(instances).hasSize(2); } + @Test + void shouldCallGetRequestOnDelegate() { + Request request = new DefaultRequest<>(new DefaultRequestContext()); + when(delegate.get()).thenReturn(Flux.just(Arrays.asList(first, second, third))); + when(delegate.get(request)).thenReturn(Flux.just(Arrays.asList(first, second))); + + List instances = supplier.get(request).blockFirst(); + + assertThat(instances).hasSize(2); + } + private DefaultServiceInstance serviceInstance(String instanceId) { return new DefaultServiceInstance(instanceId, "test", "http://test.test", 9080, false); } diff --git a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilderTests.java b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilderTests.java index 8793b8ee..7b958dd2 100644 --- a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilderTests.java +++ b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilderTests.java @@ -37,11 +37,9 @@ public class ServiceInstanceListSupplierBuilderTests { public void testBuilder() { new ApplicationContextRunner().withUserConfiguration(CacheTestConfig.class).run(context -> { ServiceInstanceListSupplier supplier = ServiceInstanceListSupplier.builder().withDiscoveryClient() - .withHealthChecks().withCaching().build(context); - assertThat(supplier).isInstanceOf(CachingServiceInstanceListSupplier.class); + .withHealthChecks().build(context); + assertThat(supplier).isInstanceOf(HealthCheckServiceInstanceListSupplier.class); DelegatingServiceInstanceListSupplier delegating = (DelegatingServiceInstanceListSupplier) supplier; - assertThat(delegating.getDelegate()).isInstanceOf(HealthCheckServiceInstanceListSupplier.class); - delegating = (DelegatingServiceInstanceListSupplier) delegating.getDelegate(); assertThat(delegating.getDelegate()).isInstanceOf(DiscoveryClientServiceInstanceListSupplier.class); }); } diff --git a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/ZonePreferenceServiceInstanceListSupplierTests.java b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/ZonePreferenceServiceInstanceListSupplierTests.java index 4fe7d39f..238cec0a 100644 --- a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/ZonePreferenceServiceInstanceListSupplierTests.java +++ b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/ZonePreferenceServiceInstanceListSupplierTests.java @@ -22,15 +22,22 @@ import java.util.HashMap; import java.util.List; import java.util.Map; +import org.junit.jupiter.api.BeforeEach; import org.junit.jupiter.api.Test; import reactor.core.publisher.Flux; import org.springframework.cloud.client.DefaultServiceInstance; import org.springframework.cloud.client.ServiceInstance; +import org.springframework.cloud.client.loadbalancer.DefaultRequest; +import org.springframework.cloud.client.loadbalancer.DefaultRequestContext; +import org.springframework.cloud.client.loadbalancer.LoadBalancerProperties; +import org.springframework.cloud.client.loadbalancer.Request; import org.springframework.cloud.loadbalancer.config.LoadBalancerZoneConfig; +import org.springframework.cloud.loadbalancer.support.LoadBalancerClientFactory; import static org.assertj.core.api.Assertions.assertThat; import static org.assertj.core.api.Assertions.assertThatCode; +import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; @@ -46,8 +53,9 @@ class ZonePreferenceServiceInstanceListSupplierTests { private final LoadBalancerZoneConfig zoneConfig = new LoadBalancerZoneConfig(null); - private final ZonePreferenceServiceInstanceListSupplier supplier = new ZonePreferenceServiceInstanceListSupplier( - delegate, zoneConfig); + private ZonePreferenceServiceInstanceListSupplier supplier; + + private final LoadBalancerClientFactory loadBalancerClientFactory = mock(LoadBalancerClientFactory.class); private final ServiceInstance first = serviceInstance("test-1", buildZoneMetadata("zone1")); @@ -59,6 +67,14 @@ class ZonePreferenceServiceInstanceListSupplierTests { private final ServiceInstance fifth = serviceInstance("test-5", buildZoneMetadata(null)); + @BeforeEach + void setUp() { + LoadBalancerProperties properties = new LoadBalancerProperties(); + properties.setCallGetWithRequestOnDelegates(true); + when(loadBalancerClientFactory.getProperties(any())).thenReturn(properties); + supplier = new ZonePreferenceServiceInstanceListSupplier(delegate, zoneConfig, loadBalancerClientFactory); + } + @Test void shouldFilterInstancesByZone() { zoneConfig.setZone("zone1"); @@ -73,6 +89,19 @@ class ZonePreferenceServiceInstanceListSupplierTests { assertThat(filtered).doesNotContain(fifth); } + @Test + void shouldCallGetRequestOnDelegate() { + zoneConfig.setZone("zone1"); + Request request = new DefaultRequest<>(new DefaultRequestContext()); + when(delegate.get()).thenReturn(Flux.just(Arrays.asList(first, second, third, fourth, fifth))); + when(delegate.get(request)).thenReturn(Flux.just(Arrays.asList(first, third, fourth, fifth))); + + List filtered = supplier.get(request).blockFirst(); + + assertThat(filtered).hasSize(1); + assertThat(filtered).containsOnly(first); + } + @Test void shouldReturnAllInstancesIfNoZoneInstances() { zoneConfig.setZone("zone1"); From d36062b8bb51b6cf39e1fa3065966f1d0c33f3bf Mon Sep 17 00:00:00 2001 From: Olga MaciaszekSharma Date: Wed, 28 Jun 2023 13:46:49 +0200 Subject: [PATCH 2/3] Update copyright and fix docs. --- docs/src/main/asciidoc/spring-cloud-commons.adoc | 2 +- .../cloud/client/loadbalancer/LoadBalancerProperties.java | 2 +- .../core/DelegatingServiceInstanceListSupplier.java | 2 +- .../core/SameInstancePreferenceServiceInstanceListSupplier.java | 2 +- .../loadbalancer/core/ServiceInstanceListSupplierBuilder.java | 2 +- .../core/ZonePreferenceServiceInstanceListSupplier.java | 2 +- .../SameInstancePreferenceServiceInstanceListSupplierTests.java | 2 +- .../core/ServiceInstanceListSupplierBuilderTests.java | 2 +- .../core/ZonePreferenceServiceInstanceListSupplierTests.java | 2 +- 9 files changed, 9 insertions(+), 9 deletions(-) diff --git a/docs/src/main/asciidoc/spring-cloud-commons.adoc b/docs/src/main/asciidoc/spring-cloud-commons.adoc index ca20aa33..893341f9 100644 --- a/docs/src/main/asciidoc/spring-cloud-commons.adoc +++ b/docs/src/main/asciidoc/spring-cloud-commons.adoc @@ -1319,7 +1319,7 @@ The per-client configuration properties work for most of the properties, apart f NOTE: For the properties where maps where already used, where you can specify a different value per-client without using the `clients` keyword (for example, `hints`, `health-check.path`), we have kept that behaviour in order to keep the library backwards compatible. It will be modified in the next major release. -NOTE: Starting with `3.1.7` in `2021.0.x` release train, `4.0.4` in `2022.0.x` release train and `4.1.0` in the `2023.0.x` release train, we have introduced the `callGetWithRequestOnDelegates` flag in `LoadBalancerProperties`. If this flag is set to `true`, `ServiceInstanceListSupplier#get(Request request)` method will be implemented to call `delegate.get(request)` in classes assignable from `DelegatingServiceInstanceListSupplier` that don't already implement that method, with the exclusion of `CachingServiceInstanceListSupplier` and `HealthCheckServiceInstanceListSupplier`, which should be placed in the instance supplier hierarchy directly after the supplier performing instance retrieval over the network, before any request-based filtering is done. For `3.1.x` and `4.0.x` the flag is set to `false` by default, and since `4.1.0` it's going to be set to `true` by default. +NOTE: Starting with `3.1.7`, we have introduced the `callGetWithRequestOnDelegates` flag in `LoadBalancerProperties`. If this flag is set to `true`, `ServiceInstanceListSupplier#get(Request request)` method will be implemented to call `delegate.get(request)` in classes assignable from `DelegatingServiceInstanceListSupplier` that don't already implement that method, with the exclusion of `CachingServiceInstanceListSupplier` and `HealthCheckServiceInstanceListSupplier`, which should be placed in the instance supplier hierarchy directly after the supplier performing instance retrieval over the network, before any request-based filtering is done. For `3.1.x` the flag is set to `false` by default, however, since `4.1.0` it's going to be set to `true` by default. == Spring Cloud Circuit Breaker diff --git a/spring-cloud-commons/src/main/java/org/springframework/cloud/client/loadbalancer/LoadBalancerProperties.java b/spring-cloud-commons/src/main/java/org/springframework/cloud/client/loadbalancer/LoadBalancerProperties.java index 5537dde1..d600e081 100644 --- a/spring-cloud-commons/src/main/java/org/springframework/cloud/client/loadbalancer/LoadBalancerProperties.java +++ b/spring-cloud-commons/src/main/java/org/springframework/cloud/client/loadbalancer/LoadBalancerProperties.java @@ -1,5 +1,5 @@ /* - * Copyright 2012-2020 the original author or authors. + * Copyright 2012-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DelegatingServiceInstanceListSupplier.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DelegatingServiceInstanceListSupplier.java index af791aa4..aaf17218 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DelegatingServiceInstanceListSupplier.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DelegatingServiceInstanceListSupplier.java @@ -1,5 +1,5 @@ /* - * Copyright 2012-2019 the original author or authors. + * Copyright 2012-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplier.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplier.java index 8317d9d9..bf504fad 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplier.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplier.java @@ -1,5 +1,5 @@ /* - * Copyright 2012-2020 the original author or authors. + * Copyright 2012-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilder.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilder.java index 0168d98b..0dcd4095 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilder.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilder.java @@ -1,5 +1,5 @@ /* - * Copyright 2013-2020 the original author or authors. + * Copyright 2013-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ZonePreferenceServiceInstanceListSupplier.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ZonePreferenceServiceInstanceListSupplier.java index 26f1c7da..dafbbad7 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ZonePreferenceServiceInstanceListSupplier.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ZonePreferenceServiceInstanceListSupplier.java @@ -1,5 +1,5 @@ /* - * Copyright 2012-2020 the original author or authors. + * Copyright 2012-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. diff --git a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplierTests.java b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplierTests.java index bf6e260b..2563db25 100644 --- a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplierTests.java +++ b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplierTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2012-2020 the original author or authors. + * Copyright 2012-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. diff --git a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilderTests.java b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilderTests.java index 7b958dd2..ecb86c02 100644 --- a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilderTests.java +++ b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilderTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2013-2020 the original author or authors. + * Copyright 2013-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. diff --git a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/ZonePreferenceServiceInstanceListSupplierTests.java b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/ZonePreferenceServiceInstanceListSupplierTests.java index 238cec0a..f04c9ba8 100644 --- a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/ZonePreferenceServiceInstanceListSupplierTests.java +++ b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/ZonePreferenceServiceInstanceListSupplierTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2012-2020 the original author or authors. + * Copyright 2012-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. From e4e56e93cfdb007df0f089647ae40ce6864d30b3 Mon Sep 17 00:00:00 2001 From: Olga MaciaszekSharma Date: Wed, 28 Jun 2023 14:30:15 +0200 Subject: [PATCH 3/3] Call `get(Request request)` on delegates in WeightedServiceInstanceListSupplier. --- .../LoadBalancerClientConfiguration.java | 4 +-- .../ServiceInstanceListSupplierBuilder.java | 13 ++++++-- .../WeightedServiceInstanceListSupplier.java | 31 +++++++++++++++-- .../LoadBalancerClientConfigurationTests.java | 8 ++--- ...ghtedServiceInstanceListSupplierTests.java | 33 ++++++++++++++++++- 5 files changed, 76 insertions(+), 13 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 ea9cd772..f0aae50c 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 @@ -139,7 +139,7 @@ public class LoadBalancerClientConfiguration { @ConditionalOnMissingBean @Conditional(WeightedConfigurationCondition.class) public ServiceInstanceListSupplier weightedServiceInstanceListSupplier(ConfigurableApplicationContext context) { - return ServiceInstanceListSupplier.builder().withDiscoveryClient().withWeighted().withCaching() + return ServiceInstanceListSupplier.builder().withDiscoveryClient().withCaching().withWeighted() .build(context); } @@ -204,7 +204,7 @@ public class LoadBalancerClientConfiguration { @ConditionalOnMissingBean @Conditional(WeightedConfigurationCondition.class) public ServiceInstanceListSupplier weightedServiceInstanceListSupplier(ConfigurableApplicationContext context) { - return ServiceInstanceListSupplier.builder().withBlockingDiscoveryClient().withWeighted().withCaching() + return ServiceInstanceListSupplier.builder().withBlockingDiscoveryClient().withCaching().withWeighted() .build(context); } diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilder.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilder.java index 7906638a..74ee2e38 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilder.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/ServiceInstanceListSupplierBuilder.java @@ -116,7 +116,11 @@ public final class ServiceInstanceListSupplierBuilder { * @return the {@link ServiceInstanceListSupplierBuilder} object */ public ServiceInstanceListSupplierBuilder withWeighted() { - DelegateCreator creator = (context, delegate) -> new WeightedServiceInstanceListSupplier(delegate); + DelegateCreator creator = (context, delegate) -> { + ReactiveLoadBalancer.Factory loadBalancerClientFactory = context + .getBean(LoadBalancerClientFactory.class); + return new WeightedServiceInstanceListSupplier(delegate, loadBalancerClientFactory); + }; this.creators.add(creator); return this; } @@ -129,8 +133,11 @@ public final class ServiceInstanceListSupplierBuilder { * @return the {@link ServiceInstanceListSupplierBuilder} object */ public ServiceInstanceListSupplierBuilder withWeighted(WeightFunction weightFunction) { - DelegateCreator creator = (context, delegate) -> new WeightedServiceInstanceListSupplier(delegate, - weightFunction); + DelegateCreator creator = (context, delegate) -> { + ReactiveLoadBalancer.Factory loadBalancerClientFactory = context + .getBean(LoadBalancerClientFactory.class); + return new WeightedServiceInstanceListSupplier(delegate, weightFunction, loadBalancerClientFactory); + }; this.creators.add(creator); return this; } diff --git a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/WeightedServiceInstanceListSupplier.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/WeightedServiceInstanceListSupplier.java index 527e6173..172038b5 100644 --- a/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/WeightedServiceInstanceListSupplier.java +++ b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/WeightedServiceInstanceListSupplier.java @@ -1,5 +1,5 @@ /* - * Copyright 2012-2022 the original author or authors. + * Copyright 2012-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -24,12 +24,15 @@ import org.apache.commons.logging.LogFactory; import reactor.core.publisher.Flux; import org.springframework.cloud.client.ServiceInstance; +import org.springframework.cloud.client.loadbalancer.Request; +import org.springframework.cloud.client.loadbalancer.reactive.ReactiveLoadBalancer; /** * A {@link ServiceInstanceListSupplier} implementation that uses weights to expand the * instances provided by delegate. * * @author Zhuozhi Ji + * @author Olga Maciaszek-Sharma */ public class WeightedServiceInstanceListSupplier extends DelegatingServiceInstanceListSupplier { @@ -41,9 +44,10 @@ public class WeightedServiceInstanceListSupplier extends DelegatingServiceInstan private final WeightFunction weightFunction; + private boolean callGetWithRequestOnDelegates; + public WeightedServiceInstanceListSupplier(ServiceInstanceListSupplier delegate) { - super(delegate); - this.weightFunction = WeightedServiceInstanceListSupplier::metadataWeightFunction; + this(delegate, WeightedServiceInstanceListSupplier::metadataWeightFunction); } public WeightedServiceInstanceListSupplier(ServiceInstanceListSupplier delegate, WeightFunction weightFunction) { @@ -51,11 +55,32 @@ public class WeightedServiceInstanceListSupplier extends DelegatingServiceInstan this.weightFunction = weightFunction; } + public WeightedServiceInstanceListSupplier(ServiceInstanceListSupplier delegate, + ReactiveLoadBalancer.Factory loadBalancerClientFactory) { + this(delegate, WeightedServiceInstanceListSupplier::metadataWeightFunction, loadBalancerClientFactory); + } + + public WeightedServiceInstanceListSupplier(ServiceInstanceListSupplier delegate, WeightFunction weightFunction, + ReactiveLoadBalancer.Factory loadBalancerClientFactory) { + super(delegate); + this.weightFunction = weightFunction; + callGetWithRequestOnDelegates = loadBalancerClientFactory.getProperties(getServiceId()) + .isCallGetWithRequestOnDelegates(); + } + @Override public Flux> get() { return delegate.get().map(this::expandByWeight); } + @Override + public Flux> get(Request request) { + if (callGetWithRequestOnDelegates) { + return delegate.get(request).map(this::expandByWeight); + } + return get(); + } + private List expandByWeight(List instances) { if (instances.size() == 0) { return instances; diff --git a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/annotation/LoadBalancerClientConfigurationTests.java b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/annotation/LoadBalancerClientConfigurationTests.java index 858d5ce4..eda9dc1c 100644 --- a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/annotation/LoadBalancerClientConfigurationTests.java +++ b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/annotation/LoadBalancerClientConfigurationTests.java @@ -120,10 +120,10 @@ class LoadBalancerClientConfigurationTests { reactiveDiscoveryClientRunner.withUserConfiguration(TestConfig.class) .withPropertyValues("spring.cloud.loadbalancer.configurations=weighted").run(context -> { ServiceInstanceListSupplier supplier = context.getBean(ServiceInstanceListSupplier.class); - then(supplier).isInstanceOf(CachingServiceInstanceListSupplier.class); + then(supplier).isInstanceOf(WeightedServiceInstanceListSupplier.class); ServiceInstanceListSupplier delegate = ((DelegatingServiceInstanceListSupplier) supplier) .getDelegate(); - then(delegate).isInstanceOf(WeightedServiceInstanceListSupplier.class); + then(delegate).isInstanceOf(CachingServiceInstanceListSupplier.class); ServiceInstanceListSupplier secondDelegate = ((DelegatingServiceInstanceListSupplier) delegate) .getDelegate(); then(secondDelegate).isInstanceOf(DiscoveryClientServiceInstanceListSupplier.class); @@ -207,10 +207,10 @@ class LoadBalancerClientConfigurationTests { blockingDiscoveryClientRunner.withUserConfiguration(RestTemplateTestConfig.class) .withPropertyValues("spring.cloud.loadbalancer.configurations=weighted").run(context -> { ServiceInstanceListSupplier supplier = context.getBean(ServiceInstanceListSupplier.class); - then(supplier).isInstanceOf(CachingServiceInstanceListSupplier.class); + then(supplier).isInstanceOf(WeightedServiceInstanceListSupplier.class); ServiceInstanceListSupplier delegate = ((DelegatingServiceInstanceListSupplier) supplier) .getDelegate(); - then(delegate).isInstanceOf(WeightedServiceInstanceListSupplier.class); + then(delegate).isInstanceOf(CachingServiceInstanceListSupplier.class); then(((DelegatingServiceInstanceListSupplier) delegate).getDelegate()) .isInstanceOf(DiscoveryClientServiceInstanceListSupplier.class); }); diff --git a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/WeightedServiceInstanceListSupplierTests.java b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/WeightedServiceInstanceListSupplierTests.java index 9bd925e2..58c00682 100644 --- a/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/WeightedServiceInstanceListSupplierTests.java +++ b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/WeightedServiceInstanceListSupplierTests.java @@ -1,5 +1,5 @@ /* - * Copyright 2012-2022 the original author or authors. + * Copyright 2012-2023 the original author or authors. * * Licensed under the Apache License, Version 2.0 (the "License"); * you may not use this file except in compliance with the License. @@ -29,9 +29,15 @@ import reactor.core.publisher.Flux; import org.springframework.cloud.client.DefaultServiceInstance; import org.springframework.cloud.client.ServiceInstance; +import org.springframework.cloud.client.loadbalancer.DefaultRequest; +import org.springframework.cloud.client.loadbalancer.DefaultRequestContext; +import org.springframework.cloud.client.loadbalancer.LoadBalancerProperties; +import org.springframework.cloud.client.loadbalancer.Request; +import org.springframework.cloud.loadbalancer.support.LoadBalancerClientFactory; import static java.util.stream.Collectors.summingInt; import static org.assertj.core.api.Assertions.assertThat; +import static org.mockito.ArgumentMatchers.any; import static org.mockito.Mockito.mock; import static org.mockito.Mockito.when; import static org.springframework.cloud.loadbalancer.core.WeightedServiceInstanceListSupplier.DEFAULT_WEIGHT; @@ -40,6 +46,7 @@ import static org.springframework.cloud.loadbalancer.core.WeightedServiceInstanc * Tests for {@link WeightedServiceInstanceListSupplier}. * * @author Zhuozhi Ji + * @author Olga Maciaszek-Sharma */ class WeightedServiceInstanceListSupplierTests { @@ -177,6 +184,30 @@ class WeightedServiceInstanceListSupplierTests { assertThat(counter).containsEntry("test-3", DEFAULT_WEIGHT); } + @Test + void shouldCallGetRequestOnDelegate() { + LoadBalancerClientFactory loadBalancerClientFactory = mock(LoadBalancerClientFactory.class); + LoadBalancerProperties properties = new LoadBalancerProperties(); + properties.setCallGetWithRequestOnDelegates(true); + when(loadBalancerClientFactory.getProperties(any())).thenReturn(properties); + ServiceInstance one = serviceInstance("test-1", Collections.emptyMap()); + ServiceInstance two = serviceInstance("test-2", Collections.emptyMap()); + ServiceInstance three = serviceInstance("test-3", buildWeightMetadata(3)); + Request request = new DefaultRequest<>(new DefaultRequestContext()); + + when(delegate.get()).thenReturn(Flux.just(Arrays.asList(one, two, three))); + when(delegate.get(request)).thenReturn(Flux.just(Arrays.asList(one, two))); + WeightedServiceInstanceListSupplier supplier = new WeightedServiceInstanceListSupplier(delegate, + loadBalancerClientFactory); + + List serviceInstances = Objects.requireNonNull(supplier.get(request).blockFirst()); + Map counter = serviceInstances.stream() + .collect(Collectors.groupingBy(ServiceInstance::getInstanceId, summingInt(e -> 1))); + assertThat(counter).containsEntry("test-1", DEFAULT_WEIGHT); + assertThat(counter).containsEntry("test-2", DEFAULT_WEIGHT); + assertThat(counter).doesNotContainEntry("test-3", 3); + } + private ServiceInstance serviceInstance(String instanceId, Map metadata) { return new DefaultServiceInstance(instanceId, "test", "localhost", 8080, false, metadata); }