使用Flux.create创建字符串流遇阻,寻求非block()循环订阅方案
问题描述
我需要通过多次调用API接口创建Flux,API返回的JSON payload包含done字段:当还有数据可获取时为false,数据耗尽时为true,此时需完成Flux。为简化问题我写了模拟代码,但运行测试时仅输出一个字符串后就挂起直至报错。我知道需要循环获取元素直到满足完成条件,但不想用block(),该怎么实现?
模拟代码
RecordProvider接口
public interface RecordProvider { Mono<String> fetchString(); }
服务类方法
private final AtomicInteger stringCounter = new AtomicInteger(0); public Flux<String> assembleStringsFlux() { return Flux.<String>create(this::generateStringMono); } private void generateStringMono(FluxSink<String> fluxSink) { recordProvider.fetchString() .subscribeOn(Schedulers.boundedElastic()) .subscribe( stringResult -> { if (stringCounter.get() == 5) { fluxSink.complete(); } else { fluxSink.next(stringResult); } }, fluxSink::error); }
测试代码
@Test void givenString() { // given // when when(recordProvider.fetchString()).thenReturn(Mono.just("bla")); var result = recordProcessor.assembleStringsFlux() .map(someString -> { System.out.println("whaaaat = " + someString); return someString; }).subscribeOn(Schedulers.boundedElastic()); // then StepVerifier.create(result) .expectNextMatches(stringSignal -> stringSignal.equals("bla")) .expectNextMatches(stringSignal -> stringSignal.equals("bla")) .expectNextMatches(stringSignal -> stringSignal.equals("bla")) .expectNextMatches(stringSignal -> stringSignal.equals("bla")) .verifyComplete(); }
问题分析
原代码的核心问题是generateStringMono仅调用了一次fetchString(),发送一个元素后没有触发后续的API调用,导致Flux只发出一个元素就停滞了,无法继续生成后续元素。
解决方案
方式1:使用Flux.expand(推荐)
expand是专门用于递归获取数据的响应式操作符,它会将每个元素映射为新的Publisher,并将所有结果合并到同一个Flux中,直到返回空Mono终止循环。
修改服务类的assembleStringsFlux方法:
private final AtomicInteger stringCounter = new AtomicInteger(0); public Flux<String> assembleStringsFlux() { // 初始调用一次API获取第一个元素 return recordProvider.fetchString() .subscribeOn(Schedulers.boundedElastic()) // 递归扩展逻辑:每次获取元素后判断是否继续调用 .expand(stringResult -> { int currentCount = stringCounter.incrementAndGet(); if (currentCount >= 5) { // 达到终止条件,返回空Mono结束循环 return Mono.empty(); } // 继续调用API获取下一个元素 return recordProvider.fetchString().subscribeOn(Schedulers.boundedElastic()); }); }
方式2:使用Flux.generate手动控制循环
generate允许手动管理元素生成逻辑,通过状态变量控制循环次数:
private final AtomicInteger stringCounter = new AtomicInteger(0); public Flux<String> assembleStringsFlux() { return Flux.generate( // 初始化状态:用计数器作为状态变量 () -> stringCounter, // 生成逻辑:异步调用API,发送元素并更新状态 (counter, sink) -> { if (counter.get() >= 5) { sink.complete(); return counter; } recordProvider.fetchString() .subscribeOn(Schedulers.boundedElastic()) .subscribe( stringResult -> { sink.next(stringResult); counter.incrementAndGet(); }, sink::error // 传递API调用错误 ); return counter; } ); }
验证效果
修改后的代码会连续调用API获取指定次数的元素,满足终止条件后自动完成Flux。原测试用的StepVerifier可以正确验证多个元素的输出,不会出现挂起或报错的情况。
内容的提问来源于stack exchange,提问作者Dusko
相关产品推荐
相关产品推荐

