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

如何使用StepVerifier验证无限Flux在5秒内至少生成2个事件?

用StepVerifier验证Flux在5秒内至少生成2个事件

你的示例Flux由两个流合并而成:一个是随机间隔发射的randomIntervalEmitter,另一个是每5秒固定发射的regularDummyUpdate。要验证5秒内至少产生2个事件,核心是利用StepVerifier的时间控制与元素断言能力,下面是具体实现方案:

测试代码实现

import org.junit.jupiter.api.Test;
import reactor.core.publisher.Flux;
import reactor.test.StepVerifier;
import java.time.Duration;
import java.util.ArrayList;
import java.util.Random;
import java.util.concurrent.TimeUnit;
import java.util.concurrent.atomic.AtomicLong;

public class FluxRateTest {

    // 提取你的示例Flux生成逻辑为可复用方法
    public Flux<String> createTestFlux() {
        final AtomicLong counter = new AtomicLong(0);
        final Random rnd = new Random();

        final Flux<String> randomIntervalEmitter = Flux.generate(generator -> {
            try {
                final long counterDivided = counter.getAndIncrement() % 12; // 修正:原代码未递增counter,这里补充递增逻辑
                if (counterDivided > 0) {
                    TimeUnit.SECONDS.sleep(rnd.nextInt(1, 10));
                } else {
                    TimeUnit.MILLISECONDS.sleep(rnd.nextInt(1, 50));
                }
                generator.next("asdf " + counterDivided);
            } catch (InterruptedException e) {
                Thread.currentThread().interrupt(); // 正确传递中断状态,替代原有的printStackTrace
            }
        });

        final Flux<String> regularDummyUpdate = Flux.interval(Duration.ofSeconds(5))
                .map(e -> "" + (88 + (System.currentTimeMillis() % 104)));
        return randomIntervalEmitter.mergeWith(regularDummyUpdate);
    }

    @Test
    void testAtLeastTwoEventsInFiveSeconds() {
        Flux<String> testFlux = createTestFlux();

        StepVerifier.create(testFlux)
                .recordWith(ArrayList::new) // 记录这段时间内发射的所有元素
                .thenAwait(Duration.ofSeconds(5)) // 等待5秒,模拟实际运行时长
                .consumeRecordedWith(events -> {
                    // 断言5秒内生成的事件数量至少为2
                    assert events.size() >= 2 : "5秒内生成的事件数量不足2个,实际为" + events.size();
                })
                .thenCancel() // 取消订阅,避免无限流导致测试挂起
                .verify();
    }
}

关键细节说明

  1. 修正原Flux逻辑:原代码中counter仅做取模操作未递增,导致counterDivided永远为0,流会一直以毫秒级间隔发射元素。测试代码中改为counter.getAndIncrement() % 12,让发射间隔符合预期的随机变化逻辑。
  2. 时间控制:thenAwait(Duration.ofSeconds(5))让StepVerifier等待指定时长,模拟真实运行场景下的时间流逝。
  3. 元素断言:通过recordWith收集这段时间内的所有元素,再用consumeRecordedWith直接断言集合大小≥2,精准满足“至少2个事件”的验证需求。
  4. 取消订阅:由于目标Flux是无限流,测试结束后必须调用thenCancel()终止订阅,否则测试会一直挂起等待新元素。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.25 17:27:25