Merge branch '1.0.x'
This commit is contained in:
@@ -3,6 +3,12 @@ include::_attributes.adoc[]
|
||||
|
||||
*{spring-cloud-version}*
|
||||
|
||||
## Usage Documentation
|
||||
|
||||
The Spring Cloud CircuitBreaker project contains implementations for Resilience4J and Spring Retry.
|
||||
The APIs implemented in Spring Cloud CircuitBreaker live in Spring Cloud Commons. The usage documentation
|
||||
for these APIs are located in the https://docs.spring.io/spring-cloud-commons/docs/current/reference/html/#spring-cloud-circuit-breake[Spring Cloud Commons documentation].
|
||||
|
||||
include::spring-cloud-circuitbreaker-resilience4j.adoc[]
|
||||
|
||||
include::spring-cloud-circuitbreaker-spring-retry.adoc[]
|
||||
|
||||
@@ -8,7 +8,6 @@
|
||||
<artifactId>spring-cloud-dependencies-parent</artifactId>
|
||||
<groupId>org.springframework.cloud</groupId>
|
||||
<version>3.0.2</version>
|
||||
<relativePath/>
|
||||
</parent>
|
||||
|
||||
<artifactId>spring-cloud-circuitbreaker-dependencies</artifactId>
|
||||
|
||||
@@ -38,8 +38,12 @@ import org.springframework.cloud.client.circuitbreaker.Customizer;
|
||||
*/
|
||||
public class Resilience4JCircuitBreaker implements CircuitBreaker {
|
||||
|
||||
static final String CIRCUIT_BREAKER_GROUP_TAG = "group";
|
||||
|
||||
private final String id;
|
||||
|
||||
private final String groupName;
|
||||
|
||||
private Resilience4jBulkheadProvider bulkheadProvider;
|
||||
|
||||
private final io.github.resilience4j.circuitbreaker.CircuitBreakerConfig circuitBreakerConfig;
|
||||
@@ -61,6 +65,7 @@ public class Resilience4JCircuitBreaker implements CircuitBreaker {
|
||||
ExecutorService executorService,
|
||||
Optional<Customizer<io.github.resilience4j.circuitbreaker.CircuitBreaker>> circuitBreakerCustomizer) {
|
||||
this.id = id;
|
||||
this.groupName = id;
|
||||
this.circuitBreakerConfig = circuitBreakerConfig;
|
||||
this.registry = circuitBreakerRegistry;
|
||||
this.timeLimiterRegistry = TimeLimiterRegistry.ofDefaults();
|
||||
@@ -69,13 +74,25 @@ public class Resilience4JCircuitBreaker implements CircuitBreaker {
|
||||
this.circuitBreakerCustomizer = circuitBreakerCustomizer;
|
||||
}
|
||||
|
||||
@Deprecated
|
||||
public Resilience4JCircuitBreaker(String id,
|
||||
io.github.resilience4j.circuitbreaker.CircuitBreakerConfig circuitBreakerConfig,
|
||||
TimeLimiterConfig timeLimiterConfig, CircuitBreakerRegistry circuitBreakerRegistry,
|
||||
TimeLimiterRegistry timeLimiterRegistry, ExecutorService executorService,
|
||||
Optional<Customizer<io.github.resilience4j.circuitbreaker.CircuitBreaker>> circuitBreakerCustomizer,
|
||||
Resilience4jBulkheadProvider bulkheadProvider) {
|
||||
this(id, id, circuitBreakerConfig, timeLimiterConfig, circuitBreakerRegistry, timeLimiterRegistry,
|
||||
executorService, circuitBreakerCustomizer, bulkheadProvider);
|
||||
}
|
||||
|
||||
public Resilience4JCircuitBreaker(String id, String groupName,
|
||||
io.github.resilience4j.circuitbreaker.CircuitBreakerConfig circuitBreakerConfig,
|
||||
TimeLimiterConfig timeLimiterConfig, CircuitBreakerRegistry circuitBreakerRegistry,
|
||||
TimeLimiterRegistry timeLimiterRegistry, ExecutorService executorService,
|
||||
Optional<Customizer<io.github.resilience4j.circuitbreaker.CircuitBreaker>> circuitBreakerCustomizer,
|
||||
Resilience4jBulkheadProvider bulkheadProvider) {
|
||||
this.id = id;
|
||||
this.groupName = groupName;
|
||||
this.circuitBreakerConfig = circuitBreakerConfig;
|
||||
this.registry = circuitBreakerRegistry;
|
||||
this.timeLimiterRegistry = timeLimiterRegistry;
|
||||
@@ -87,16 +104,18 @@ public class Resilience4JCircuitBreaker implements CircuitBreaker {
|
||||
|
||||
@Override
|
||||
public <T> T run(Supplier<T> toRun, Function<Throwable, T> fallback) {
|
||||
TimeLimiter timeLimiter = timeLimiterRegistry.timeLimiter(id, timeLimiterConfig);
|
||||
final io.vavr.collection.Map<String, String> tags = io.vavr.collection.HashMap.of(CIRCUIT_BREAKER_GROUP_TAG,
|
||||
this.groupName);
|
||||
TimeLimiter timeLimiter = timeLimiterRegistry.timeLimiter(id, timeLimiterConfig, tags);
|
||||
Supplier<Future<T>> futureSupplier = () -> executorService.submit(toRun::get);
|
||||
Callable restrictedCall = TimeLimiter.decorateFutureSupplier(timeLimiter, futureSupplier);
|
||||
|
||||
io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry.circuitBreaker(id,
|
||||
circuitBreakerConfig);
|
||||
Callable restrictedCall = TimeLimiter.decorateFutureSupplier(timeLimiter, futureSupplier);
|
||||
io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry.circuitBreaker(this.id,
|
||||
this.circuitBreakerConfig, tags);
|
||||
circuitBreakerCustomizer.ifPresent(customizer -> customizer.customize(defaultCircuitBreaker));
|
||||
|
||||
if (bulkheadProvider != null) {
|
||||
return bulkheadProvider.run(id, toRun, fallback, defaultCircuitBreaker, timeLimiter);
|
||||
return bulkheadProvider.run(this.groupName, toRun, fallback, defaultCircuitBreaker, timeLimiter, tags);
|
||||
}
|
||||
else {
|
||||
Callable<T> callable = io.github.resilience4j.circuitbreaker.CircuitBreaker
|
||||
|
||||
@@ -19,6 +19,7 @@ package org.springframework.cloud.circuitbreaker.resilience4j;
|
||||
import java.util.HashMap;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.ConcurrentHashMap;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.function.Function;
|
||||
@@ -48,6 +49,8 @@ public class Resilience4JCircuitBreakerFactory extends
|
||||
|
||||
private ExecutorService executorService = Executors.newCachedThreadPool();
|
||||
|
||||
private ConcurrentHashMap<String, ExecutorService> executorServices = new ConcurrentHashMap<>();
|
||||
|
||||
private Map<String, Customizer<CircuitBreaker>> circuitBreakerCustomizers = new HashMap<>();
|
||||
|
||||
@Deprecated
|
||||
@@ -101,11 +104,16 @@ public class Resilience4JCircuitBreakerFactory extends
|
||||
@Override
|
||||
public Resilience4JCircuitBreaker create(String id) {
|
||||
Assert.hasText(id, "A CircuitBreaker must have an id.");
|
||||
Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config = getConfigurations()
|
||||
.computeIfAbsent(id, defaultConfiguration);
|
||||
return new Resilience4JCircuitBreaker(id, config.getCircuitBreakerConfig(), config.getTimeLimiterConfig(),
|
||||
circuitBreakerRegistry, timeLimiterRegistry, executorService,
|
||||
Optional.ofNullable(circuitBreakerCustomizers.get(id)), bulkheadProvider);
|
||||
return create(id, id, this.executorService);
|
||||
}
|
||||
|
||||
@Override
|
||||
public Resilience4JCircuitBreaker create(String id, String groupName) {
|
||||
Assert.hasText(id, "A CircuitBreaker must have an id.");
|
||||
Assert.hasText(groupName, "A CircuitBreaker must have a group name.");
|
||||
final ExecutorService groupExecutorService = executorServices.computeIfAbsent(groupName,
|
||||
group -> Executors.newCachedThreadPool());
|
||||
return create(id, groupName, groupExecutorService);
|
||||
}
|
||||
|
||||
public void addCircuitBreakerCustomizer(Customizer<CircuitBreaker> customizer, String... ids) {
|
||||
@@ -114,4 +122,14 @@ public class Resilience4JCircuitBreakerFactory extends
|
||||
}
|
||||
}
|
||||
|
||||
private Resilience4JCircuitBreaker create(String id, String groupName,
|
||||
ExecutorService circuitBreakerExecutorService) {
|
||||
Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config = getConfigurations()
|
||||
.computeIfAbsent(id, defaultConfiguration);
|
||||
return new Resilience4JCircuitBreaker(id, groupName, config.getCircuitBreakerConfig(),
|
||||
config.getTimeLimiterConfig(), circuitBreakerRegistry, timeLimiterRegistry,
|
||||
circuitBreakerExecutorService, Optional.ofNullable(circuitBreakerCustomizers.get(id)),
|
||||
bulkheadProvider);
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -99,25 +99,26 @@ public class Resilience4jBulkheadProvider {
|
||||
}
|
||||
|
||||
public <T> T run(String id, Supplier<T> toRun, Function<Throwable, T> fallback, CircuitBreaker circuitBreaker,
|
||||
TimeLimiter timeLimiter) {
|
||||
Supplier<CompletionStage<T>> bulkheadCall = decorateBulkhead(id, toRun);
|
||||
TimeLimiter timeLimiter, io.vavr.collection.Map<String, String> tags) {
|
||||
Supplier<CompletionStage<T>> bulkheadCall = decorateBulkhead(id, tags, toRun);
|
||||
final Callable<T> timeLimiterCall = decorateTimeLimiter(bulkheadCall, timeLimiter);
|
||||
final Callable<T> circuitBreakerCall = circuitBreaker.decorateCallable(timeLimiterCall);
|
||||
return Try.of(circuitBreakerCall::call).recover(fallback).get();
|
||||
}
|
||||
|
||||
private <T> Supplier<CompletionStage<T>> decorateBulkhead(final String id, final Supplier<T> supplier) {
|
||||
private <T> Supplier<CompletionStage<T>> decorateBulkhead(final String id,
|
||||
final io.vavr.collection.Map<String, String> tags, final Supplier<T> supplier) {
|
||||
Resilience4jBulkheadConfigurationBuilder.BulkheadConfiguration configuration = configurations
|
||||
.computeIfAbsent(id, defaultConfiguration);
|
||||
|
||||
if (bulkheadRegistry.find(id).isPresent() && !threadPoolBulkheadRegistry.find(id).isPresent()) {
|
||||
Bulkhead bulkhead = bulkheadRegistry.bulkhead(id, configuration.getBulkheadConfig());
|
||||
Bulkhead bulkhead = bulkheadRegistry.bulkhead(id, configuration.getBulkheadConfig(), tags);
|
||||
CompletableFuture<T> asyncCall = CompletableFuture.supplyAsync(supplier);
|
||||
return Bulkhead.decorateCompletionStage(bulkhead, () -> asyncCall);
|
||||
}
|
||||
else {
|
||||
ThreadPoolBulkhead threadPoolBulkhead = threadPoolBulkheadRegistry.bulkhead(id,
|
||||
configuration.getThreadPoolBulkheadConfig());
|
||||
configuration.getThreadPoolBulkheadConfig(), tags);
|
||||
return threadPoolBulkhead.decorateSupplier(supplier);
|
||||
}
|
||||
}
|
||||
|
||||
@@ -107,6 +107,16 @@ public class Resilience4JCircuitBreakerIntegrationTest {
|
||||
assertThat(service.normal()).isEqualTo("normal");
|
||||
assertThat(((List) rest.getForObject("/actuator/metrics", Map.class).get("names"))
|
||||
.contains("resilience4j.circuitbreaker.calls")).isTrue();
|
||||
|
||||
// CircuitBreaker and TimeLimiter should have 3 metrics: name, kind, group
|
||||
assertThat(((List) rest
|
||||
.getForObject("/actuator/metrics/resilience4j.circuitbreaker.calls",
|
||||
Map.class)
|
||||
.get("availableTags"))).hasSize(3);
|
||||
assertThat(((List) rest
|
||||
.getForObject("/actuator/metrics/resilience4j.timelimiter.calls",
|
||||
Map.class)
|
||||
.get("availableTags"))).hasSize(3);
|
||||
}
|
||||
|
||||
@Configuration(proxyBeanMethods = false)
|
||||
|
||||
@@ -39,6 +39,15 @@ public class Resilience4JCircuitBreakerTest {
|
||||
assertThat(cb.run(() -> "foobar")).isEqualTo("foobar");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void runWithGroupName() {
|
||||
CircuitBreaker cb = new Resilience4JCircuitBreakerFactory(
|
||||
CircuitBreakerRegistry.ofDefaults(), TimeLimiterRegistry.ofDefaults(),
|
||||
null).create("foo", "groupFoo");
|
||||
assertThat(cb.run(() -> "foobar")).isEqualTo("foobar");
|
||||
|
||||
}
|
||||
|
||||
@Test
|
||||
public void runWithFallback() {
|
||||
CircuitBreaker cb = new Resilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(),
|
||||
@@ -48,6 +57,16 @@ public class Resilience4JCircuitBreakerTest {
|
||||
}, t -> "fallback")).isEqualTo("fallback");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void runWithFallbackAndGroupName() {
|
||||
CircuitBreaker cb = new Resilience4JCircuitBreakerFactory(
|
||||
CircuitBreakerRegistry.ofDefaults(), TimeLimiterRegistry.ofDefaults(),
|
||||
null).create("foo", "groupFoo");
|
||||
assertThat((String) cb.run(() -> {
|
||||
throw new RuntimeException("boom");
|
||||
}, t -> "fallback")).isEqualTo("fallback");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void runWithBulkheadProvider() {
|
||||
CircuitBreaker cb = new Resilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(),
|
||||
@@ -56,6 +75,15 @@ public class Resilience4JCircuitBreakerTest {
|
||||
assertThat(cb.run(() -> "foobar")).isEqualTo("foobar");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void runWithBulkheadProviderAndGroupName() {
|
||||
CircuitBreaker cb = new Resilience4JCircuitBreakerFactory(
|
||||
CircuitBreakerRegistry.ofDefaults(), TimeLimiterRegistry.ofDefaults(),
|
||||
new Resilience4jBulkheadProvider(ThreadPoolBulkheadRegistry.ofDefaults(),
|
||||
BulkheadRegistry.ofDefaults())).create("foo", "groupFoo");
|
||||
assertThat(cb.run(() -> "foobar")).isEqualTo("foobar");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void runWithFallbackBulkheadProvider() {
|
||||
CircuitBreaker cb = new Resilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(),
|
||||
@@ -66,4 +94,15 @@ public class Resilience4JCircuitBreakerTest {
|
||||
}, t -> "fallback")).isEqualTo("fallback");
|
||||
}
|
||||
|
||||
@Test
|
||||
public void runWithFallbackBulkheadProviderAndGroupName() {
|
||||
CircuitBreaker cb = new Resilience4JCircuitBreakerFactory(
|
||||
CircuitBreakerRegistry.ofDefaults(), TimeLimiterRegistry.ofDefaults(),
|
||||
new Resilience4jBulkheadProvider(ThreadPoolBulkheadRegistry.ofDefaults(),
|
||||
BulkheadRegistry.ofDefaults())).create("foo", "groupFoo");
|
||||
assertThat((String) cb.run(() -> {
|
||||
throw new RuntimeException("boom");
|
||||
}, t -> "fallback")).isEqualTo("fallback");
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user