如何使用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(); } }
关键细节说明
- 修正原Flux逻辑:原代码中
counter仅做取模操作未递增,导致counterDivided永远为0,流会一直以毫秒级间隔发射元素。测试代码中改为counter.getAndIncrement() % 12,让发射间隔符合预期的随机变化逻辑。 - 时间控制:
thenAwait(Duration.ofSeconds(5))让StepVerifier等待指定时长,模拟真实运行场景下的时间流逝。 - 元素断言:通过
recordWith收集这段时间内的所有元素,再用consumeRecordedWith直接断言集合大小≥2,精准满足“至少2个事件”的验证需求。 - 取消订阅:由于目标Flux是无限流,测试结束后必须调用
thenCancel()终止订阅,否则测试会一直挂起等待新元素。
内容的提问来源于stack exchange,提问作者Lubo
相关产品推荐
相关产品推荐

