You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何使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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.18 11:47:25