Handle exceptions and timeouts new (#767)

* Handle timeouts and exceptions while retrieving instances.

* Update docs.
This commit is contained in:
Olga Maciaszek-Sharma
2020-05-27 13:11:47 -05:00
committed by GitHub
parent 87e5d7a62b
commit da23fe4f07
5 changed files with 147 additions and 16 deletions

View File

@@ -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);
}
}
}

View File

@@ -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"
}
]
}

View File

@@ -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);
}
}