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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 06:45:58