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

如何正确测试Sink(Processor)输出的Flux?

解决StepVerifier与Sink配合时错过测试信号的问题

问题根源

你使用的multicast().directBestEffort()类型Sink,在没有订阅者的情况下发送消息会直接丢弃消息。原代码里发送消息的操作早于StepVerifier的订阅(StepVerifier仅在调用verify()时才完成订阅),导致所有测试信号丢失,最终验证失败。

核心解决方案:让订阅先于消息发送

通过StepVerifier的thenRun()方法,把发送消息的逻辑放到订阅完成之后执行,确保订阅者存在时再推送消息:

import reactor.core.publisher.Sinks;
import reactor.test.StepVerifier;
import java.time.Duration;

public class TestBed {
    public static void main(String[] args) {
        class StringProcessor {
            public final Sinks.Many<String> sink = Sinks.many().multicast().directBestEffort();

            public void httpPostWebhookController(String inputData) {
                sink.emitNext(
                        inputData.toLowerCase() + " " + inputData.toUpperCase(),
                        (signalType, emitResult) -> {
                            System.out.println("error, signalType=" + signalType + "; emitResult=" + emitResult);
                            return false;
                        }
                );
            }
        }

        final StringProcessor stringProcessor = new StringProcessor();
        
        StepVerifier.create(stringProcessor.sink.asFlux())
                .expectSubscription()
                // 订阅完成后再发送测试消息
                .thenRun(() -> {
                    stringProcessor.httpPostWebhookController("asdf");
                    stringProcessor.httpPostWebhookController("Qw");
                })
                .expectNext("asdf ASDF")
                .expectNext("qw QW")
                .thenCancel()
                .verify(Duration.ofSeconds(2));
    }
}

原理说明

StepVerifier按步骤顺序推进执行:

  1. create()初始化验证器,但此时未订阅Flux
  2. expectSubscription()触发验证器完成订阅操作
  3. thenRun()执行消息发送逻辑,此时订阅者已存在,Sink能正常推送消息给StepVerifier
  4. 后续expectNext()可匹配到推送的消息,完成验证

可选方案:使用带缓存的Sink

如果业务允许晚订阅的消费者获取历史消息,可改用带缓存的Sink类型,即使先发送消息再订阅,StepVerifier也能收到之前的信号:

// 修改Sink创建方式,启用缓存
public final Sinks.Many<String> sink = Sinks.many().multicast().onBackpressureBuffer();

这种方式适合消息量不大的场景,需注意缓存带来的内存占用问题。

内容的提问来源于stack exchange,提问作者Lubo

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.22 08:33:35