diff --git a/docs/src/main/asciidoc/spring-cloud-circuitbreaker.adoc b/docs/src/main/asciidoc/spring-cloud-circuitbreaker.adoc index 6223eeb..0047220 100755 --- a/docs/src/main/asciidoc/spring-cloud-circuitbreaker.adoc +++ b/docs/src/main/asciidoc/spring-cloud-circuitbreaker.adoc @@ -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[] diff --git a/spring-cloud-circuitbreaker-dependencies/pom.xml b/spring-cloud-circuitbreaker-dependencies/pom.xml index 7abbb28..430b6e5 100644 --- a/spring-cloud-circuitbreaker-dependencies/pom.xml +++ b/spring-cloud-circuitbreaker-dependencies/pom.xml @@ -8,7 +8,6 @@ spring-cloud-dependencies-parent org.springframework.cloud 3.0.2 - spring-cloud-circuitbreaker-dependencies diff --git a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreaker.java b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreaker.java index 60dcc18..2889730 100644 --- a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreaker.java +++ b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreaker.java @@ -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> 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> 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> 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 run(Supplier toRun, Function fallback) { - TimeLimiter timeLimiter = timeLimiterRegistry.timeLimiter(id, timeLimiterConfig); + final io.vavr.collection.Map tags = io.vavr.collection.HashMap.of(CIRCUIT_BREAKER_GROUP_TAG, + this.groupName); + TimeLimiter timeLimiter = timeLimiterRegistry.timeLimiter(id, timeLimiterConfig, tags); Supplier> 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 callable = io.github.resilience4j.circuitbreaker.CircuitBreaker diff --git a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerFactory.java b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerFactory.java index 86a88a9..a3ed12e 100644 --- a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerFactory.java +++ b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerFactory.java @@ -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 executorServices = new ConcurrentHashMap<>(); + private Map> 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 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); + } + } diff --git a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4jBulkheadProvider.java b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4jBulkheadProvider.java index 95460a8..d267d1d 100644 --- a/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4jBulkheadProvider.java +++ b/spring-cloud-circuitbreaker-resilience4j/src/main/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4jBulkheadProvider.java @@ -99,25 +99,26 @@ public class Resilience4jBulkheadProvider { } public T run(String id, Supplier toRun, Function fallback, CircuitBreaker circuitBreaker, - TimeLimiter timeLimiter) { - Supplier> bulkheadCall = decorateBulkhead(id, toRun); + TimeLimiter timeLimiter, io.vavr.collection.Map tags) { + Supplier> bulkheadCall = decorateBulkhead(id, tags, toRun); final Callable timeLimiterCall = decorateTimeLimiter(bulkheadCall, timeLimiter); final Callable circuitBreakerCall = circuitBreaker.decorateCallable(timeLimiterCall); return Try.of(circuitBreakerCall::call).recover(fallback).get(); } - private Supplier> decorateBulkhead(final String id, final Supplier supplier) { + private Supplier> decorateBulkhead(final String id, + final io.vavr.collection.Map tags, final Supplier 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 asyncCall = CompletableFuture.supplyAsync(supplier); return Bulkhead.decorateCompletionStage(bulkhead, () -> asyncCall); } else { ThreadPoolBulkhead threadPoolBulkhead = threadPoolBulkheadRegistry.bulkhead(id, - configuration.getThreadPoolBulkheadConfig()); + configuration.getThreadPoolBulkheadConfig(), tags); return threadPoolBulkhead.decorateSupplier(supplier); } } diff --git a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerIntegrationTest.java b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerIntegrationTest.java index 062bdf7..6dde096 100644 --- a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerIntegrationTest.java +++ b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerIntegrationTest.java @@ -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) diff --git a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerTest.java b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerTest.java index edd5a57..ddcbc27 100644 --- a/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerTest.java +++ b/spring-cloud-circuitbreaker-resilience4j/src/test/java/org/springframework/cloud/circuitbreaker/resilience4j/Resilience4JCircuitBreakerTest.java @@ -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"); + } + }