Add support for reactive Resilience4jBulkheadProvider to provide bulkhead support for reactive operations (Mono and Flux).

This commit is contained in:
Yavor Chamov
2024-12-22 14:31:22 +02:00
parent d3c6af8d63
commit 93dc27b588
4 changed files with 16 additions and 7 deletions

View File

@@ -30,8 +30,6 @@ import org.springframework.context.ConfigurableApplicationContext;
import static org.assertj.core.api.Assertions.assertThat;
/**
* @author Ryan Baxter
* @author Thomas Vitale
* @author Yavor Chamov
*/
@RunWith(ModifiedClassPathRunner.class)

View File

@@ -37,6 +37,9 @@ import org.springframework.stereotype.Service;
import org.springframework.test.annotation.DirtiesContext;
import org.springframework.test.context.junit4.SpringRunner;
/**
* @author Yavor Chamov
*/
@RunWith(SpringRunner.class)
@SpringBootTest(classes = ReactiveResilience4JBulkheadAndTimeLimiterIntegrationTest.Application.class)
@DirtiesContext

View File

@@ -44,7 +44,8 @@ public class ReactiveResilience4JCircuitBreakerTest {
@Test
public void runMono() {
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(),
TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()), new Resilience4JConfigurationProperties())
TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()),
new Resilience4JConfigurationProperties())
.create("foo");
assertThat(Mono.just("foobar").transform(cb::run).block()).isEqualTo("foobar");
}
@@ -52,7 +53,8 @@ public class ReactiveResilience4JCircuitBreakerTest {
@Test
public void runMonoWithFallback() {
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(),
TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()), new Resilience4JConfigurationProperties())
TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()),
new Resilience4JConfigurationProperties())
.create("foo");
assertThat(Mono.error(new RuntimeException("boom"))
.transform(it -> cb.run(it, t -> Mono.just("fallback")))
@@ -62,7 +64,8 @@ public class ReactiveResilience4JCircuitBreakerTest {
@Test
public void runFlux() {
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(),
TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()), new Resilience4JConfigurationProperties())
TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()),
new Resilience4JConfigurationProperties())
.create("foo");
assertThat(Flux.just("foobar", "hello world").transform(cb::run).collectList().block())
.isEqualTo(Arrays.asList("foobar", "hello world"));
@@ -71,7 +74,8 @@ public class ReactiveResilience4JCircuitBreakerTest {
@Test
public void runFluxWithFallback() {
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(),
TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()), new Resilience4JConfigurationProperties())
TimeLimiterRegistry.ofDefaults(), new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()),
new Resilience4JConfigurationProperties())
.create("foo");
assertThat(Flux.error(new RuntimeException("boom"))
.transform(it -> cb.run(it, t -> Flux.just("fallback")))
@@ -87,7 +91,8 @@ public class ReactiveResilience4JCircuitBreakerTest {
public void runWithDefaultTimeLimiter() {
final TimeLimiterRegistry timeLimiterRegistry = TimeLimiterRegistry.ofDefaults();
ReactiveCircuitBreaker cb = new ReactiveResilience4JCircuitBreakerFactory(CircuitBreakerRegistry.ofDefaults(),
timeLimiterRegistry, new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()), new Resilience4JConfigurationProperties())
timeLimiterRegistry, new ReactiveResilience4jBulkheadProvider(BulkheadRegistry.ofDefaults()),
new Resilience4JConfigurationProperties())
.create("foo");
assertThat(Mono.fromCallable(() -> {

View File

@@ -54,6 +54,9 @@ import org.springframework.web.bind.annotation.RestController;
import static org.assertj.core.api.AssertionsForInterfaceTypes.assertThat;
import static org.springframework.boot.test.context.SpringBootTest.WebEnvironment.RANDOM_PORT;
/**
* @author Yavor Chamov
*/
@RunWith(SpringRunner.class)
@SpringBootTest(webEnvironment = RANDOM_PORT, classes = ReactiveResilience4jBulkheadIntegrationTest.Application.class,
properties = {"management.endpoints.web.exposure.include=*"})