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

控制Flux.generate发射值实现Reactor动态DAG顺序任务执行

问题根源

  • Flux.generate默认遵循Reactor的预取规则,会提前生成最多队列容量个元素放入缓冲,即便下游concatMap是顺序执行,上游依然会提前下发任务,导致任务提前触发doOnNext打印,和实际执行顺序脱节
  • takeUntil放在concatMap下游,触发终止时会直接取消上游,不会等待当前正在执行的任务完成
  • 外部上下文ctx跨线程修改存在可见性问题,可能导致生成下一个任务时读取到旧状态

正确实现方案

推荐方案:使用Mono.expand实现动态迭代

Mono.expand是专门用来实现「依赖前序结果生成下一个异步任务」场景的操作符,天然等待前一个任务完成后才触发下一个任务的生成,完全匹配你要的do{}while逻辑,实现如下:

// 启动第一个任务
Mono<OrchestrationContext> firstTask = getFirstTaskFromDAG().run(ctx);
firstTask
    // 前序任务执行完成后,根据返回的上下文生成下一个任务
    .expand(prevCtx -> {
        ReactiveTask<OrchestrationContext> nextTask = deriveNextStep(prevCtx.getLastExecutedStep(), prevCtx.getDecisionData());
        if ("END".equals(nextTask.getName()) || !evaluateStatus(prevCtx)) {
            // 返回空表示终止迭代
            return Mono.empty();
        }
        return nextTask.run(prevCtx);
    })
    .publishOn(this.factory.getSharedSchedulerPool())
    .onErrorResume(throwable -> buildResponse(ctx, throwable))
    .doOnCancel(() -> log.info("Task cancelled"))
    .doOnComplete(() -> log.info("Completed flow"))
    .subscribe();

若要保留Flux.generate的调整方案

需要强制关闭预取机制,保证上游每次只生成1个任务,并且调整操作符顺序:

Flux.generate(
        (SynchronousSink<ReactiveTask<OrchestrationContext>> synchronousSink) -> {
            ReactiveTask<OrchestrationContext> task = null;
            if (ctx.getLastExecutedStep() == null) {
                task = getFirstTaskFromDAG();
            } else {
                task = deriveNextStep(ctx.getLastExecutedStep(), ctx.getDecisionData());
            }
            if ("END".equals(task.getName())) {
                synchronousSink.complete();
                return;
            }
            synchronousSink.next(task);
        })
    // 指定publishOn预取数为1,每次只向上游请求1个元素
    .publishOn(this.factory.getSharedSchedulerPool(), 1)
    .doOnNext(task -> log.info("On next => {}", task.getTaskName()))
    // concatMap同样指定预取数为1,不提前缓冲任务
    .concatMap(reactiveTask -> reactiveTask.run(ctx), 1)
    .takeUntil(context -> !evaluateStatus(context, tasks))
    .onErrorResume(throwable -> buildResponse(ctx, throwable))
    .doOnCancel(() -> log.info("Task cancelled"))
    .doOnComplete(() -> log.info("Completed flow"))
    .subscribe();

额外注意点

  • 上下文ctx建议每次任务执行后返回新的不可变实例,避免多线程修改的可见性问题
  • 任务执行逻辑里的状态修改需要做好线程安全控制,避免deriveNextStep读取到中间状态

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 07:06:03