diff --git a/docs/src/main/asciidoc/_configprops.adoc b/docs/src/main/asciidoc/_configprops.adoc index df7fa55c..7185f92c 100644 --- a/docs/src/main/asciidoc/_configprops.adoc +++ b/docs/src/main/asciidoc/_configprops.adoc @@ -37,6 +37,7 @@ |spring.cloud.loadbalancer.cache.capacity | `+++256+++` | Initial cache capacity expressed as int. |spring.cloud.loadbalancer.cache.enabled | `+++true+++` | Enables Spring Cloud LoadBalancer caching mechanism. |spring.cloud.loadbalancer.cache.ttl | `+++35s+++` | Time To Live - time counted from writing of the record, after which cache entries are expired, expressed as a {@link Duration}. The property {@link String} has to be in keeping with the appropriate syntax as specified in Spring Boot StringToDurationConverter. @see StringToDurationConverter.java +|spring.cloud.loadbalancer.call-get-with-request-on-delegates | `+++false+++` | 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 |spring.cloud.loadbalancer.clients | | |spring.cloud.loadbalancer.configurations | `+++default+++` | Enables a predefined LoadBalancer configuration. |spring.cloud.loadbalancer.eager-load.clients | | Names of the clients. diff --git a/docs/src/main/asciidoc/spring-cloud-commons.adoc b/docs/src/main/asciidoc/spring-cloud-commons.adoc index 0a5dc7ff..58cc2808 100644 --- a/docs/src/main/asciidoc/spring-cloud-commons.adoc +++ b/docs/src/main/asciidoc/spring-cloud-commons.adoc @@ -947,6 +947,12 @@ 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. + +==== + === Weighted Load-Balancing To enable weighted load-balancing, we provide the `WeightedServiceInstanceListSupplier`. We use `WeightFunction` to calculate the weight of each instance. @@ -1011,7 +1017,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: @@ -1025,8 +1031,8 @@ public class CustomLoadBalancerConfiguration { ConfigurableApplicationContext context) { return ServiceInstanceListSupplier.builder() .withDiscoveryClient() + .withCaching() .withZonePreference() - .withCaching() .build(context); } } @@ -1089,6 +1095,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. @@ -1173,8 +1185,8 @@ public class CustomLoadBalancerConfiguration { ConfigurableApplicationContext context) { return ServiceInstanceListSupplier.builder() .withDiscoveryClient() + .withCaching() .withHints() - .withCaching() .build(context); } } @@ -1284,11 +1296,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 @@ -1356,6 +1375,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 `4.0.4`, 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 `4.0.x` the flag is set to `false` by default, however, since `4.1.0` it's going to be set to `true` by default. + === AOT and Native Image Support Since `4.0.0`, Spring Cloud LoadBalancer supports Spring AOT transformations and native images. However, to use this feature, you need to explicitly define your `LoadBalancerClient` service IDs. You can do so by using the `value` or `name` attributes of the `@LoadBalancerClient` annotation or as values of the `spring.cloud.loadbalancer.eager-load.clients` property. 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 c7ab68fd..0e2c2121 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. @@ -72,6 +72,19 @@ public class LoadBalancerProperties { */ private StickySession stickySession = new StickySession(); + /** + * 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; } @@ -125,6 +138,36 @@ public class LoadBalancerProperties { return xForwarded; } + /** + * 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 8f71163b..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 @@ -94,7 +94,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); } @@ -120,8 +120,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 @@ -130,8 +130,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); } @Bean @@ -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); } @@ -165,8 +165,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 @@ -185,8 +185,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 @@ -195,8 +195,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); } @Bean @@ -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/DelegatingServiceInstanceListSupplier.java b/spring-cloud-loadbalancer/src/main/java/org/springframework/cloud/loadbalancer/core/DelegatingServiceInstanceListSupplier.java index 80b6600f..42bf25fb 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 @@ -40,12 +40,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 8d8e63f2..bd63ecbc 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 @@ -40,10 +42,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(); @@ -54,6 +65,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 c6aa0f85..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 @@ -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. @@ -58,8 +58,6 @@ public final class ServiceInstanceListSupplierBuilder { private Creator baseCreator; - private DelegateCreator cachingCreator; - private final List creators = new ArrayList<>(); ServiceInstanceListSupplierBuilder() { @@ -118,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; } @@ -131,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; } @@ -174,8 +179,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; } @@ -217,8 +224,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; @@ -232,8 +240,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; @@ -254,19 +263,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) { @@ -277,6 +282,7 @@ public final class ServiceInstanceListSupplierBuilder { } return delegate; }; + creators.add(creator); return this; } @@ -323,9 +329,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/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/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..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. @@ -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 3e26296e..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 @@ -93,10 +93,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); @@ -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); @@ -136,10 +136,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); @@ -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/SameInstancePreferenceServiceInstanceListSupplierTests.java b/spring-cloud-loadbalancer/src/test/java/org/springframework/cloud/loadbalancer/core/SameInstancePreferenceServiceInstanceListSupplierTests.java index e7b58b2f..81995b89 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,11 +19,17 @@ 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; @@ -43,8 +49,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"); @@ -52,6 +59,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))); @@ -73,7 +88,7 @@ class SameInstancePreferenceServiceInstanceListSupplierTests { } @Test - void shouldReturnAllInstancesFromDelegateIfPreviouslySelectedInstanceIfAvailable() { + void shouldReturnAllInstancesFromDelegateIfPreviouslySelectedInstanceIsNotAvailable() { when(delegate.get()).thenReturn(Flux.just(Arrays.asList(second, third))); supplier.selectedServiceInstance(first); @@ -82,6 +97,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); + } + @Test void shouldCallSelectedServiceInstanceOnItsDelegate() { ServiceInstance firstInstance = serviceInstance("test-4"); 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 66458777..9bfb0f38 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. @@ -37,11 +37,9 @@ public class ServiceInstanceListSupplierBuilderTests { public void testBuilder() { new ApplicationContextRunner().withUserConfiguration(CacheTestConfig.class).run(context -> { ServiceInstanceListSupplier supplier = ServiceInstanceListSupplier.builder().withDiscoveryClient() - .withHealthChecks().withWeighted().withCaching().build(context); - assertThat(supplier).isInstanceOf(CachingServiceInstanceListSupplier.class); + .withHealthChecks().withWeighted().build(context); + assertThat(supplier).isInstanceOf(WeightedServiceInstanceListSupplier.class); DelegatingServiceInstanceListSupplier delegating = (DelegatingServiceInstanceListSupplier) supplier; - assertThat(delegating.getDelegate()).isInstanceOf(WeightedServiceInstanceListSupplier.class); - delegating = (DelegatingServiceInstanceListSupplier) delegating.getDelegate(); 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/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); } 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..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. @@ -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");