Simplifying customization of resilience4j circuit breakers
This commit is contained in:
@@ -72,21 +72,17 @@ public class ReactiveHystrixCircuitBreakerIntegrationTest {
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Customizer<ReactiveCircuitBreakerFactory<HystrixObservableCommand.Setter,
|
||||
ReactiveHystrixCircuitBreakerFactory.ReactiveHystrixConfigBuilder>> customizer() {
|
||||
public Customizer<ReactiveHystrixCircuitBreakerFactory> customizer() {
|
||||
return factory -> factory.configure(builder -> builder.commandProperties(
|
||||
HystrixCommandProperties.Setter().withExecutionTimeoutInMilliseconds(2000)), "slow");
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Customizer<ReactiveCircuitBreakerFactory<HystrixObservableCommand.Setter,
|
||||
ReactiveHystrixCircuitBreakerFactory.ReactiveHystrixConfigBuilder>> defaultConfig() {
|
||||
return factory -> factory.configureDefault(id -> {
|
||||
return HystrixObservableCommand.Setter.withGroupKey(HystrixCommandGroupKey.Factory.asKey(id))
|
||||
.andCommandPropertiesDefaults(HystrixCommandProperties.Setter()
|
||||
.withExecutionTimeoutInMilliseconds(4000));
|
||||
|
||||
});
|
||||
public Customizer<ReactiveHystrixCircuitBreakerFactory> defaultConfig() {
|
||||
return factory -> factory.configureDefault(id -> HystrixObservableCommand.Setter
|
||||
.withGroupKey(HystrixCommandGroupKey.Factory.asKey(id))
|
||||
.andCommandPropertiesDefaults(HystrixCommandProperties.Setter()
|
||||
.withExecutionTimeoutInMilliseconds(4000)));
|
||||
}
|
||||
|
||||
@Service
|
||||
|
||||
@@ -37,16 +37,12 @@ import org.springframework.context.annotation.Configuration;
|
||||
@ConditionalOnClass(name = {"reactor.core.publisher.Mono", "reactor.core.publisher.Flux"})
|
||||
public class ReactiveResilience4JAutoConfiguration {
|
||||
|
||||
@Autowired(required = false)
|
||||
public List<Customizer<CircuitBreaker>> circuitBreakerCustomizers = new ArrayList<>();
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean(ReactiveCircuitBreakerFactory.class)
|
||||
public ReactiveCircuitBreakerFactory reactiveResilience4JCircuitBreakerFactory() {
|
||||
return new ReactiveResilience4JCircuitBreakerFactory(circuitBreakerCustomizers);
|
||||
return new ReactiveResilience4JCircuitBreakerFactory();
|
||||
}
|
||||
|
||||
|
||||
@Configuration
|
||||
@ConditionalOnClass(name = {"reactor.core.publisher.Mono", "reactor.core.publisher.Flux"})
|
||||
public static class ReactiveResilience4JCustomizerConfiguration {
|
||||
|
||||
@@ -22,6 +22,7 @@ import reactor.core.publisher.Flux;
|
||||
import reactor.core.publisher.Mono;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.TimeoutException;
|
||||
import java.util.function.Function;
|
||||
|
||||
@@ -37,21 +38,21 @@ public class ReactiveResilience4JCircuitBreaker implements ReactiveCircuitBreake
|
||||
private String id;
|
||||
private Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config;
|
||||
private CircuitBreakerRegistry registry;
|
||||
private List<Customizer<CircuitBreaker>> circuitBreakerCustomizers;
|
||||
private Optional<Customizer<CircuitBreaker>> circuitBreakerCustomizer;
|
||||
|
||||
public ReactiveResilience4JCircuitBreaker(String id, Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config,
|
||||
CircuitBreakerRegistry circuitBreakerRegistry,
|
||||
List<Customizer<CircuitBreaker>> circuitBreakerCustomizers) {
|
||||
Optional<Customizer<CircuitBreaker>> circuitBreakerCustomizer) {
|
||||
this.id = id;
|
||||
this.config = config;
|
||||
this.registry = circuitBreakerRegistry;
|
||||
this.circuitBreakerCustomizers = circuitBreakerCustomizers;
|
||||
this.circuitBreakerCustomizer = circuitBreakerCustomizer;
|
||||
}
|
||||
|
||||
@Override
|
||||
public <T> Mono<T> run(Mono<T> toRun, Function<Throwable, Mono<T>> fallback) {
|
||||
io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry.circuitBreaker(id, config.getCircuitBreakerConfig());
|
||||
circuitBreakerCustomizers.forEach(circuitBreakerCustomizer -> circuitBreakerCustomizer.customize(defaultCircuitBreaker));
|
||||
circuitBreakerCustomizer.ifPresent(customizer -> customizer.customize(defaultCircuitBreaker));
|
||||
Mono<T> toReturn = toRun.transform(CircuitBreakerOperator.of(defaultCircuitBreaker))
|
||||
.timeout(config.getTimeLimiterConfig().getTimeoutDuration())
|
||||
// Since we are using the Mono timeout we need to tell the circuit breaker about the error
|
||||
@@ -64,7 +65,7 @@ public class ReactiveResilience4JCircuitBreaker implements ReactiveCircuitBreake
|
||||
|
||||
public<T> Flux<T> run(Flux<T> toRun, Function<Throwable, Flux<T>> fallback) {
|
||||
io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry.circuitBreaker(id, config.getCircuitBreakerConfig());
|
||||
circuitBreakerCustomizers.forEach(circuitBreakerCustomizer -> circuitBreakerCustomizer.customize(defaultCircuitBreaker));
|
||||
circuitBreakerCustomizer.ifPresent(customizer -> customizer.customize(defaultCircuitBreaker));
|
||||
Flux<T> toReturn = toRun.transform(CircuitBreakerOperator.of(defaultCircuitBreaker))
|
||||
.timeout(config.getTimeLimiterConfig().getTimeoutDuration())
|
||||
// Since we are using the Flux timeout we need to tell the circuit breaker about the error
|
||||
|
||||
@@ -20,7 +20,10 @@ import io.github.resilience4j.circuitbreaker.CircuitBreakerConfig;
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry;
|
||||
import io.github.resilience4j.timelimiter.TimeLimiterConfig;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.function.Function;
|
||||
|
||||
import org.springframework.cloud.circuitbreaker.commons.Customizer;
|
||||
@@ -40,18 +43,14 @@ public class ReactiveResilience4JCircuitBreakerFactory extends ReactiveCircuitBr
|
||||
.build();
|
||||
|
||||
private CircuitBreakerRegistry circuitBreakerRegistry = CircuitBreakerRegistry.ofDefaults();
|
||||
private List<Customizer<CircuitBreaker>> circuitBreakerCustomizers;
|
||||
|
||||
public ReactiveResilience4JCircuitBreakerFactory(List<Customizer<CircuitBreaker>> circuitBreakerCustomizers) {
|
||||
this.circuitBreakerCustomizers = circuitBreakerCustomizers;
|
||||
}
|
||||
private Map<String, Customizer<CircuitBreaker>> circuitBreakerCustomizers = new HashMap<>();
|
||||
|
||||
@Override
|
||||
public ReactiveCircuitBreaker create(String id) {
|
||||
Assert.hasText(id, "A CircuitBreaker must have an id.");
|
||||
Resilience4JConfigBuilder.Resilience4JCircuitBreakerConfiguration config = getConfigurations().computeIfAbsent(id, defaultConfiguration);
|
||||
return new ReactiveResilience4JCircuitBreaker(id, config,
|
||||
circuitBreakerRegistry, circuitBreakerCustomizers);
|
||||
circuitBreakerRegistry, Optional.ofNullable(circuitBreakerCustomizers.get(id)));
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -67,4 +66,10 @@ public class ReactiveResilience4JCircuitBreakerFactory extends ReactiveCircuitBr
|
||||
public void configureCircuitBreakerRegistry(CircuitBreakerRegistry registry) {
|
||||
this.circuitBreakerRegistry = registry;
|
||||
}
|
||||
|
||||
public void addCircuitBreakerCustomizer(Customizer<CircuitBreaker> customizer, String... ids) {
|
||||
for(String id : ids) {
|
||||
circuitBreakerCustomizers.put(id, customizer);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -36,13 +36,10 @@ import org.springframework.context.annotation.Configuration;
|
||||
@Configuration
|
||||
public class Resilience4JAutoConfiguration {
|
||||
|
||||
@Autowired(required = false)
|
||||
public List<Customizer<CircuitBreaker>> circuitBreakerCustomizers = new ArrayList<>();
|
||||
|
||||
@Bean
|
||||
@ConditionalOnMissingBean(CircuitBreakerFactory.class)
|
||||
public CircuitBreakerFactory resilience4jCircuitBreakerFactory() {
|
||||
return new Resilience4JCircuitBreakerFactory(circuitBreakerCustomizers);
|
||||
return new Resilience4JCircuitBreakerFactory();
|
||||
}
|
||||
|
||||
@Configuration
|
||||
|
||||
@@ -22,6 +22,7 @@ import io.github.resilience4j.timelimiter.TimeLimiterConfig;
|
||||
import io.vavr.control.Try;
|
||||
|
||||
import java.util.List;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.Callable;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Future;
|
||||
@@ -41,17 +42,17 @@ public class Resilience4JCircuitBreaker implements CircuitBreaker {
|
||||
private CircuitBreakerRegistry registry;
|
||||
private TimeLimiterConfig timeLimiterConfig;
|
||||
private ExecutorService executorService;
|
||||
private List<Customizer<io.github.resilience4j.circuitbreaker.CircuitBreaker>> circuitBreakerCustomizers;
|
||||
private Optional<Customizer<io.github.resilience4j.circuitbreaker.CircuitBreaker>> circuitBreakerCustomizer;
|
||||
|
||||
public Resilience4JCircuitBreaker(String id, io.github.resilience4j.circuitbreaker.CircuitBreakerConfig circuitBreakerConfig,
|
||||
TimeLimiterConfig timeLimiterConfig, CircuitBreakerRegistry circuitBreakerRegistry,
|
||||
ExecutorService executorService, List<Customizer<io.github.resilience4j.circuitbreaker.CircuitBreaker>> circuitBreakerCustomizers) {
|
||||
ExecutorService executorService, Optional<Customizer<io.github.resilience4j.circuitbreaker.CircuitBreaker>> circuitBreakerCustomizer) {
|
||||
this.id = id;
|
||||
this.circuitBreakerConfig = circuitBreakerConfig;
|
||||
this.registry = circuitBreakerRegistry;
|
||||
this.timeLimiterConfig = timeLimiterConfig;
|
||||
this.executorService = executorService;
|
||||
this.circuitBreakerCustomizers = circuitBreakerCustomizers;
|
||||
this.circuitBreakerCustomizer = circuitBreakerCustomizer;
|
||||
}
|
||||
|
||||
@Override
|
||||
@@ -62,7 +63,7 @@ public class Resilience4JCircuitBreaker implements CircuitBreaker {
|
||||
.decorateFutureSupplier(timeLimiter, futureSupplier);
|
||||
|
||||
io.github.resilience4j.circuitbreaker.CircuitBreaker defaultCircuitBreaker = registry.circuitBreaker(id, circuitBreakerConfig);
|
||||
circuitBreakerCustomizers.forEach(circuitBreakerCustomizer -> circuitBreakerCustomizer.customize(defaultCircuitBreaker));
|
||||
circuitBreakerCustomizer.ifPresent(customizer -> customizer.customize(defaultCircuitBreaker));
|
||||
Callable<T> callable = io.github.resilience4j.circuitbreaker.CircuitBreaker
|
||||
.decorateCallable(defaultCircuitBreaker, restrictedCall);
|
||||
return Try.of(callable::call).recover(fallback).get();
|
||||
|
||||
@@ -21,7 +21,10 @@ import io.github.resilience4j.circuitbreaker.CircuitBreakerConfig;
|
||||
import io.github.resilience4j.circuitbreaker.CircuitBreakerRegistry;
|
||||
import io.github.resilience4j.timelimiter.TimeLimiterConfig;
|
||||
|
||||
import java.util.HashMap;
|
||||
import java.util.List;
|
||||
import java.util.Map;
|
||||
import java.util.Optional;
|
||||
import java.util.concurrent.ExecutorService;
|
||||
import java.util.concurrent.Executors;
|
||||
import java.util.function.Function;
|
||||
@@ -42,11 +45,7 @@ public class Resilience4JCircuitBreakerFactory extends CircuitBreakerFactory<Res
|
||||
|
||||
private CircuitBreakerRegistry circuitBreakerRegistry = CircuitBreakerRegistry.ofDefaults();
|
||||
private ExecutorService executorService = Executors.newSingleThreadExecutor();
|
||||
private List<Customizer<CircuitBreaker>> circuitBreakerCustomizers;
|
||||
|
||||
public Resilience4JCircuitBreakerFactory(List<Customizer<CircuitBreaker>> circuitBreakerCustomizes) {
|
||||
this.circuitBreakerCustomizers = circuitBreakerCustomizes;
|
||||
}
|
||||
private Map<String, Customizer<CircuitBreaker>> circuitBreakerCustomizers = new HashMap<>();
|
||||
|
||||
@Override
|
||||
protected Resilience4JConfigBuilder configBuilder(String id) {
|
||||
@@ -71,7 +70,13 @@ public class Resilience4JCircuitBreakerFactory extends CircuitBreakerFactory<Res
|
||||
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, executorService, circuitBreakerCustomizers);
|
||||
circuitBreakerRegistry, executorService, Optional.ofNullable(circuitBreakerCustomizers.get(id)));
|
||||
}
|
||||
|
||||
public void addCircuitBreakerCustomizer(Customizer<CircuitBreaker> customizer, String... ids) {
|
||||
for(String id: ids) {
|
||||
circuitBreakerCustomizers.put(id, customizer);
|
||||
}
|
||||
}
|
||||
|
||||
}
|
||||
|
||||
@@ -106,53 +106,22 @@ public class ReactiveResilience4JCircuitBreakerIntegrationTest {
|
||||
|
||||
@Bean
|
||||
public Customizer<ReactiveResilience4JCircuitBreakerFactory> slowCusomtizer() {
|
||||
return factory -> factory.configure(builder -> builder
|
||||
.timeLimiterConfig(TimeLimiterConfig.custom().timeoutDuration(Duration.ofSeconds(2)).build())
|
||||
.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults()), "slow", "slowflux");
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Customizer<ReactiveResilience4JCircuitBreakerFactory> defaultCustomizer() {
|
||||
return factory -> factory.configureDefault(id -> new Resilience4JConfigBuilder(id)
|
||||
.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
|
||||
.timeLimiterConfig(TimeLimiterConfig.custom().timeoutDuration(Duration.ofSeconds(4)).build()).build());
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Customizer<CircuitBreaker> slowCircuitBreakerCustomizer() {
|
||||
return circuitBreaker -> {
|
||||
if("slow".equals(circuitBreaker.getName())) {
|
||||
circuitBreaker.getEventPublisher().onError(slowErrorConsumer).onSuccess(slowSuccessConsumer);
|
||||
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Customizer<CircuitBreaker> normalCircuitBreakerCustomizer() {
|
||||
return circuitBreaker -> {
|
||||
if("normal".equals(circuitBreaker.getName())) {
|
||||
circuitBreaker.getEventPublisher().onError(normalErrorConsumer).onSuccess(normalSuccessConsumer);
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Customizer<CircuitBreaker> slowFluxCircuitBreakerCustomizer() {
|
||||
return circuitBreaker -> {
|
||||
if("slowflux".equals(circuitBreaker.getName())) {
|
||||
circuitBreaker.getEventPublisher().onError(slowFluxErrorConsumer).onSuccess(slowFluxSuccessConsumer);
|
||||
|
||||
}
|
||||
};
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Customizer<CircuitBreaker> normalFluxCircuitBreakerCustomizer() {
|
||||
return circuitBreaker -> {
|
||||
if("normalflux".equals(circuitBreaker.getName())) {
|
||||
circuitBreaker.getEventPublisher().onError(normalFluxErrorConsumer).onSuccess(normalFluxSuccessConsumer);
|
||||
}
|
||||
return factory -> {
|
||||
factory.configureDefault(id -> new Resilience4JConfigBuilder(id)
|
||||
.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
|
||||
.timeLimiterConfig(TimeLimiterConfig.custom().timeoutDuration(Duration.ofSeconds(4)).build()).build());
|
||||
factory.configure(builder -> builder
|
||||
.timeLimiterConfig(TimeLimiterConfig.custom().timeoutDuration(Duration.ofSeconds(2)).build())
|
||||
.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults()), "slow", "slowflux");
|
||||
factory.addCircuitBreakerCustomizer(circuitBreaker ->
|
||||
circuitBreaker.getEventPublisher().onError(slowErrorConsumer).onSuccess(slowSuccessConsumer),
|
||||
"slow" );
|
||||
factory.addCircuitBreakerCustomizer(circuitBreaker -> circuitBreaker.getEventPublisher().onError(normalErrorConsumer).onSuccess(normalSuccessConsumer),
|
||||
"normal");
|
||||
factory.addCircuitBreakerCustomizer(circuitBreaker -> circuitBreaker.getEventPublisher().onError(slowFluxErrorConsumer).onSuccess(slowFluxSuccessConsumer),
|
||||
"slowflux");
|
||||
factory.addCircuitBreakerCustomizer(circuitBreaker -> circuitBreaker.getEventPublisher().onError(normalFluxErrorConsumer).onSuccess(normalFluxSuccessConsumer),
|
||||
"normalflux");
|
||||
};
|
||||
}
|
||||
|
||||
|
||||
@@ -33,25 +33,25 @@ public class ReactiveResilience4JCircuitBreakerTest {
|
||||
|
||||
@Test
|
||||
public void runMono() {
|
||||
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(Collections.emptyList()).create("foo");
|
||||
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory().create("foo");
|
||||
assertEquals("foobar", cb.run(Mono.just("foobar")).block());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void runMonoWithFallback() {
|
||||
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(Collections.emptyList()).create("foo");
|
||||
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory().create("foo");
|
||||
assertEquals("fallback", cb.run(Mono.error(new RuntimeException("boom")), t -> Mono.just("fallback")).block());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void runFlux() {
|
||||
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(Collections.emptyList()).create("foo");
|
||||
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory().create("foo");
|
||||
assertEquals(Arrays.asList("foobar", "hello world"), cb.run(Flux.just("foobar", "hello world")).collectList().block());
|
||||
}
|
||||
|
||||
@Test
|
||||
public void runFluxWithFallback() {
|
||||
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(Collections.emptyList()).create("foo");
|
||||
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory().create("foo");
|
||||
assertEquals(Arrays.asList("fallback"), cb.run(Flux.error(new RuntimeException("boom")), t -> Flux.just("fallback")).collectList().block());
|
||||
}
|
||||
|
||||
|
||||
@@ -64,16 +64,14 @@ public class Resilience4JCircuitBreakerIntegrationTest {
|
||||
|
||||
@Bean
|
||||
public Customizer<Resilience4JCircuitBreakerFactory> slowCustomizer() {
|
||||
return factory -> factory.configure(builder -> builder.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
|
||||
.timeLimiterConfig(TimeLimiterConfig.custom().timeoutDuration(Duration.ofSeconds(2)).build()), "slow");
|
||||
}
|
||||
|
||||
@Bean
|
||||
public Customizer<Resilience4JCircuitBreakerFactory> defaultCustomizer() {
|
||||
return factory -> factory.configureDefault(id -> new Resilience4JConfigBuilder(id)
|
||||
.timeLimiterConfig(TimeLimiterConfig.custom().timeoutDuration(Duration.ofSeconds(4)).build())
|
||||
.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
|
||||
.build());
|
||||
return factory -> {
|
||||
factory.configure(builder -> builder.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
|
||||
.timeLimiterConfig(TimeLimiterConfig.custom().timeoutDuration(Duration.ofSeconds(2)).build()), "slow");
|
||||
factory.configureDefault(id -> new Resilience4JConfigBuilder(id)
|
||||
.timeLimiterConfig(TimeLimiterConfig.custom().timeoutDuration(Duration.ofSeconds(4)).build())
|
||||
.circuitBreakerConfig(CircuitBreakerConfig.ofDefaults())
|
||||
.build());
|
||||
};
|
||||
}
|
||||
|
||||
@Service
|
||||
|
||||
@@ -30,13 +30,13 @@ public class Resilience4JCircuitBreakerTest {
|
||||
|
||||
@Test
|
||||
public void run() {
|
||||
CircuitBreaker cb = new Resilience4JCircuitBreakerFactory(Collections.emptyList()).create("foo");
|
||||
CircuitBreaker cb = new Resilience4JCircuitBreakerFactory().create("foo");
|
||||
assertEquals("foobar", cb.run(() -> "foobar"));
|
||||
}
|
||||
|
||||
@Test
|
||||
public void runWithFallback() {
|
||||
CircuitBreaker cb = new Resilience4JCircuitBreakerFactory(Collections.emptyList()).create("foo");
|
||||
CircuitBreaker cb = new Resilience4JCircuitBreakerFactory().create("foo");
|
||||
assertEquals("fallback", cb.run(() -> {
|
||||
throw new RuntimeException("boom");
|
||||
}, t -> "fallback"));
|
||||
|
||||
Reference in New Issue
Block a user