Merge remote-tracking branch 'origin/2.2.x'
# Conflicts: # docs/src/main/asciidoc/_configprops.adoc
This commit is contained in:
@@ -16,11 +16,16 @@
|
||||
|
||||
package org.springframework.cloud.loadbalancer.core;
|
||||
|
||||
import java.time.Duration;
|
||||
import java.util.ArrayList;
|
||||
import java.util.List;
|
||||
|
||||
import org.apache.commons.logging.Log;
|
||||
import org.apache.commons.logging.LogFactory;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.core.scheduler.Schedulers;
|
||||
|
||||
import org.springframework.boot.convert.DurationStyle;
|
||||
import org.springframework.cloud.client.ServiceInstance;
|
||||
import org.springframework.cloud.client.discovery.DiscoveryClient;
|
||||
import org.springframework.cloud.client.discovery.ReactiveDiscoveryClient;
|
||||
@@ -39,6 +44,16 @@ import static org.springframework.cloud.loadbalancer.support.LoadBalancerClientF
|
||||
public class DiscoveryClientServiceInstanceListSupplier
|
||||
implements ServiceInstanceListSupplier {
|
||||
|
||||
/**
|
||||
* Property that establishes the timeout for calls to service discovery.
|
||||
*/
|
||||
public static final String SERVICE_DISCOVERY_TIMEOUT = "spring.cloud.loadbalancer.service-discovery.timeout";
|
||||
|
||||
private static final Log LOG = LogFactory
|
||||
.getLog(DiscoveryClientServiceInstanceListSupplier.class);
|
||||
|
||||
private Duration timeout = Duration.ofSeconds(30);
|
||||
|
||||
private final String serviceId;
|
||||
|
||||
private final Flux<List<ServiceInstance>> serviceInstances;
|
||||
@@ -46,16 +61,30 @@ public class DiscoveryClientServiceInstanceListSupplier
|
||||
public DiscoveryClientServiceInstanceListSupplier(DiscoveryClient delegate,
|
||||
Environment environment) {
|
||||
this.serviceId = environment.getProperty(PROPERTY_NAME);
|
||||
resolveTimeout(environment);
|
||||
this.serviceInstances = Flux
|
||||
.defer(() -> Flux.just(delegate.getInstances(serviceId)))
|
||||
.subscribeOn(Schedulers.boundedElastic());
|
||||
.timeout(timeout, Flux.defer(() -> {
|
||||
logTimeout();
|
||||
return Flux.just(new ArrayList<>());
|
||||
})).onErrorResume(error -> {
|
||||
logException(error);
|
||||
return Flux.just(new ArrayList<>());
|
||||
}).subscribeOn(Schedulers.boundedElastic());
|
||||
}
|
||||
|
||||
public DiscoveryClientServiceInstanceListSupplier(ReactiveDiscoveryClient delegate,
|
||||
Environment environment) {
|
||||
this.serviceId = environment.getProperty(PROPERTY_NAME);
|
||||
this.serviceInstances = Flux
|
||||
.defer(() -> delegate.getInstances(serviceId).collectList().flux());
|
||||
resolveTimeout(environment);
|
||||
this.serviceInstances = Flux.defer(() -> delegate.getInstances(serviceId)
|
||||
.collectList().flux().timeout(timeout, Flux.defer(() -> {
|
||||
logTimeout();
|
||||
return Flux.just(new ArrayList<>());
|
||||
})).onErrorResume(error -> {
|
||||
logException(error);
|
||||
return Flux.just(new ArrayList<>());
|
||||
}));
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -68,4 +97,28 @@ public class DiscoveryClientServiceInstanceListSupplier
|
||||
return serviceInstances;
|
||||
}
|
||||
|
||||
private void resolveTimeout(Environment environment) {
|
||||
String providedTimeout = environment.getProperty(SERVICE_DISCOVERY_TIMEOUT);
|
||||
if (providedTimeout != null) {
|
||||
timeout = DurationStyle.detectAndParse(providedTimeout);
|
||||
}
|
||||
}
|
||||
|
||||
private void logTimeout() {
|
||||
if (LOG.isDebugEnabled()) {
|
||||
LOG.debug(String.format(
|
||||
"Timeout occurred while retrieving instances for service %s."
|
||||
+ "The instances could not be retrieved during %s",
|
||||
serviceId, timeout));
|
||||
}
|
||||
}
|
||||
|
||||
private void logException(Throwable error) {
|
||||
if (LOG.isDebugEnabled()) {
|
||||
LOG.debug(String.format(
|
||||
"Exception occurred while retrieving instances for service %s",
|
||||
serviceId), error);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -10,6 +10,17 @@
|
||||
"name": "spring.cloud.loadbalancer.zone",
|
||||
"type": "java.lang.String",
|
||||
"description": "Spring Cloud LoadBalancer zone."
|
||||
},
|
||||
{
|
||||
"name": "spring.cloud.loadbalancer.service-discovery.timeout",
|
||||
"description": "String representation of Duration of the timeout for calls to service discovery.",
|
||||
"type": "java.lang.String"
|
||||
},
|
||||
{
|
||||
"defaultValue": true,
|
||||
"name": "spring.cloud.loadbalancer.cache.enabled",
|
||||
"description": "Enables Spring Cloud LoadBalancer caching mechanism.",
|
||||
"type": "java.lang.Boolean"
|
||||
}
|
||||
]
|
||||
}
|
||||
|
||||
@@ -16,9 +16,13 @@
|
||||
|
||||
package org.springframework.cloud.loadbalancer.core;
|
||||
|
||||
import java.time.Duration;
|
||||
|
||||
import org.assertj.core.util.Lists;
|
||||
import org.junit.jupiter.api.BeforeEach;
|
||||
import org.junit.jupiter.api.Test;
|
||||
import org.mockito.internal.stubbing.answers.AnswersWithDelay;
|
||||
import org.mockito.internal.stubbing.answers.Returns;
|
||||
import reactor.core.publisher.Flux;
|
||||
import reactor.test.StepVerifier;
|
||||
|
||||
@@ -29,6 +33,7 @@ import org.springframework.mock.env.MockEnvironment;
|
||||
|
||||
import static org.mockito.Mockito.mock;
|
||||
import static org.mockito.Mockito.when;
|
||||
import static org.springframework.cloud.loadbalancer.core.DiscoveryClientServiceInstanceListSupplier.SERVICE_DISCOVERY_TIMEOUT;
|
||||
|
||||
/**
|
||||
* Tests for {@link DiscoveryClientServiceInstanceListSupplier}.
|
||||
@@ -39,6 +44,8 @@ class DiscoveryClientServiceInstanceListSupplierTests {
|
||||
|
||||
private static final String SERVICE_ID = "test";
|
||||
|
||||
private static final Duration VERIFICATION_TIMEOUT = Duration.ofSeconds(10);
|
||||
|
||||
private final MockEnvironment environment = new MockEnvironment();
|
||||
|
||||
private final ReactiveDiscoveryClient reactiveDiscoveryClient = mock(
|
||||
@@ -66,9 +73,10 @@ class DiscoveryClientServiceInstanceListSupplierTests {
|
||||
supplier = new DiscoveryClientServiceInstanceListSupplier(
|
||||
reactiveDiscoveryClient, environment);
|
||||
return supplier.get();
|
||||
}).expectSubscription().expectNext(
|
||||
Lists.list(instance("1host", false), instance("2host-secure", true)))
|
||||
.thenCancel().verify();
|
||||
}).expectSubscription()
|
||||
.expectNext(Lists.list(instance("1host", false),
|
||||
instance("2host-secure", true)))
|
||||
.thenCancel().verify(VERIFICATION_TIMEOUT);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -81,7 +89,7 @@ class DiscoveryClientServiceInstanceListSupplierTests {
|
||||
StepVerifier.withVirtualTime(() -> supplier.get()).expectSubscription()
|
||||
.expectNext(Lists.list(instance("1host", false),
|
||||
instance("2host-secure", true)))
|
||||
.thenCancel().verify();
|
||||
.thenCancel().verify(VERIFICATION_TIMEOUT);
|
||||
|
||||
when(reactiveDiscoveryClient.getInstances(SERVICE_ID))
|
||||
.thenReturn(Flux.just(instance("1host", false),
|
||||
@@ -90,7 +98,31 @@ class DiscoveryClientServiceInstanceListSupplierTests {
|
||||
StepVerifier.withVirtualTime(() -> supplier.get()).expectSubscription()
|
||||
.expectNext(Lists.list(instance("1host", false),
|
||||
instance("2host-secure", true), instance("3host", false)))
|
||||
.thenCancel().verify();
|
||||
.thenCancel().verify(VERIFICATION_TIMEOUT);
|
||||
}
|
||||
|
||||
@Test
|
||||
void shouldReturnEmptyInstancesListOnException() {
|
||||
when(reactiveDiscoveryClient.getInstances(SERVICE_ID))
|
||||
.thenReturn(Flux.error(new RuntimeException("Exception")));
|
||||
|
||||
StepVerifier.withVirtualTime(() -> {
|
||||
supplier = new DiscoveryClientServiceInstanceListSupplier(
|
||||
reactiveDiscoveryClient, environment);
|
||||
return supplier.get();
|
||||
}).expectSubscription().expectNext(Lists.emptyList()).thenCancel()
|
||||
.verify(VERIFICATION_TIMEOUT);
|
||||
}
|
||||
|
||||
@Test
|
||||
void shouldReturnEmptyInstancesListOnTimeout() {
|
||||
environment.setProperty(SERVICE_DISCOVERY_TIMEOUT, "100ms");
|
||||
when(reactiveDiscoveryClient.getInstances(SERVICE_ID)).thenReturn(Flux.never());
|
||||
StepVerifier
|
||||
.create(new DiscoveryClientServiceInstanceListSupplier(
|
||||
reactiveDiscoveryClient, environment).get())
|
||||
.expectSubscription().expectNext(Lists.emptyList()).thenCancel()
|
||||
.verify(VERIFICATION_TIMEOUT);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -102,9 +134,10 @@ class DiscoveryClientServiceInstanceListSupplierTests {
|
||||
supplier = new DiscoveryClientServiceInstanceListSupplier(discoveryClient,
|
||||
environment);
|
||||
return supplier.get();
|
||||
}).expectSubscription().expectNext(
|
||||
Lists.list(instance("1host", false), instance("2host-secure", true)))
|
||||
.thenCancel().verify();
|
||||
}).expectSubscription()
|
||||
.expectNext(Lists.list(instance("1host", false),
|
||||
instance("2host-secure", true)))
|
||||
.thenCancel().verify(VERIFICATION_TIMEOUT);
|
||||
}
|
||||
|
||||
@Test
|
||||
@@ -123,7 +156,33 @@ class DiscoveryClientServiceInstanceListSupplierTests {
|
||||
}).expectSubscription()
|
||||
.expectNext(Lists.list(instance("1host", false),
|
||||
instance("2host-secure", true), instance("3host", false)))
|
||||
.thenCancel().verify();
|
||||
.thenCancel().verify(VERIFICATION_TIMEOUT);
|
||||
}
|
||||
|
||||
@Test
|
||||
void shouldReturnEmptyInstancesListOnExceptionBlockingClient() {
|
||||
when(discoveryClient.getInstances(SERVICE_ID))
|
||||
.thenThrow(new RuntimeException("Exception"));
|
||||
|
||||
StepVerifier.withVirtualTime(() -> {
|
||||
supplier = new DiscoveryClientServiceInstanceListSupplier(discoveryClient,
|
||||
environment);
|
||||
return supplier.get();
|
||||
}).expectSubscription().expectNext(Lists.emptyList()).thenCancel()
|
||||
.verify(VERIFICATION_TIMEOUT);
|
||||
}
|
||||
|
||||
@Test
|
||||
void shouldReturnEmptyInstancesListOnTimeoutBlockingClient() {
|
||||
environment.setProperty(SERVICE_DISCOVERY_TIMEOUT, "100ms");
|
||||
when(discoveryClient.getInstances(SERVICE_ID)).thenAnswer(
|
||||
new AnswersWithDelay(200, new Returns(Lists.list(instance("1host", false),
|
||||
instance("2host-secure", true), instance("3host", false)))));
|
||||
StepVerifier
|
||||
.create(new DiscoveryClientServiceInstanceListSupplier(discoveryClient,
|
||||
environment).get())
|
||||
.expectSubscription().expectNext(Lists.emptyList()).thenCancel()
|
||||
.verify(VERIFICATION_TIMEOUT);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user