如何使Reactor TestPublisher首次订阅返回错误、后续返回值测试重试?
解决Reactor重试机制测试:首次失败后重试成功的场景
问题的核心在于TestPublisher的信号序列是固定的——一旦你给它发送了错误信号,每次订阅都会重播这个错误,自然没法实现“第一次失败、第二次成功”的重试场景。要解决这个问题,关键是让每次订阅时能动态返回不同的信号,这里推荐用Mono.defer()或Flux.defer()来包装动态逻辑。
实现方案
假设你的目标方法结构大致如下(依赖一个外部Publisher,通过retryWhen处理重试):
public class MyService { private final Mono<MyObject> fetchPublisher; public MyService(Mono<MyObject> fetchPublisher) { this.fetchPublisher = fetchPublisher; } public Mono<MyObject> createMyObject() { return fetchPublisher .retryWhen(Retry.backoff(3, Duration.ofMillis(100))); // 示例重试策略 } } class MyObject { private final String value; public MyObject(String value) { this.value = value; } public String getValue() { return value; } }
方法1:用defer动态生成订阅信号
通过AtomicInteger计数订阅次数,第一次返回错误,后续返回正常结果:
import org.junit.jupiter.api.Test; import reactor.core.publisher.Mono; import reactor.test.StepVerifier; import reactor.util.retry.Retry; import java.time.Duration; import java.util.concurrent.atomic.AtomicInteger; import static org.assertj.core.api.Assertions.assertThat; public class MyServiceTest { @Test void createMyObject_shouldRetryAndSucceedAfterFirstFailure() { AtomicInteger attemptCounter = new AtomicInteger(0); // 每次订阅时动态决定返回错误还是成功 Mono<MyObject> mockFetch = Mono.defer(() -> { int attempt = attemptCounter.incrementAndGet(); if (attempt == 1) { return Mono.error(new RuntimeException("Initial fetch failed")); } else { return Mono.just(new MyObject("success")); } }); MyService service = new MyService(mockFetch); StepVerifier.create(service.createMyObject()) .expectNextMatches(obj -> obj.getValue().equals("success")) .verifyComplete(); // 验证总共触发了2次订阅(1次初始+1次重试) assertThat(attemptCounter.get()).isEqualTo(2); } }
方法2:配合TestPublisher(不推荐,仅作参考)
如果一定要用TestPublisher,可以通过defer配合计数器手动控制信号,但这种方式更繁琐:
@Test void createMyObject_withTestPublisher() { TestPublisher<MyObject> testPublisher = TestPublisher.create(); AtomicInteger subscribeCount = new AtomicInteger(0); Mono<MyObject> wrappedPublisher = Mono.defer(() -> { int count = subscribeCount.incrementAndGet(); if (count == 1) { testPublisher.error(new RuntimeException("First attempt failed")); } else { testPublisher.next(new MyObject("success")); testPublisher.complete(); } return testPublisher.mono(); }); MyService service = new MyService(wrappedPublisher); StepVerifier.create(service.createMyObject()) .expectNextMatches(obj -> obj.getValue().equals("success")) .verifyComplete(); assertThat(subscribeCount.get()).isEqualTo(2); }
关键原理
Mono.defer()的作用是延迟Publisher的创建直到订阅发生,所以每次订阅都会重新执行内部的逻辑。这样就能完美适配重试场景:第一次订阅触发错误,触发retryWhen的重试逻辑,第二次订阅时执行新的逻辑返回成功信号,避免了TestPublisher固定信号导致的无限重试问题。
内容的提问来源于stack exchange,提问作者Leprechaun
相关产品推荐
相关产品推荐

