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

如何用Reactive WebFlux实现Pipeline设计模式的顺序执行?

问题解决方案

你当前的execute方法存在两个核心问题:一是所有处理器都在处理同一个初始对象,没有实现结果的顺序传递;二是语法错误(.t属于无效代码)。以下是修正后的实现,确保处理器按排序后的顺序逐个执行,且每个处理器的输入是前一个处理器的输出:

修正后的PipeLine类代码

import org.slf4j.Logger;
import org.slf4j.LoggerFactory;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;
import java.util.ArrayList;
import java.util.List;
import java.util.function.Function;
import static java.util.Comparator.comparingInt;

public class PipeLine {
    private static final Logger log = LoggerFactory.getLogger(PipeLine.class);
    private final List<SeedPreProcessor<PreProcessorDocument>> allProcessors;

    public PipeLine(List<SeedPreProcessor<PreProcessorDocument>> allProcessors) {
        this.allProcessors = new ArrayList<>(allProcessors);
        this.allProcessors.sort(comparingInt(SeedPreProcessor::order));
    }

    public Mono<PreProcessorDocument> execute(String url) {
        log.info("Start processing URL = {}", url);
        PreProcessorDocument initial = new PreProcessorDocument(url);

        // 按顺序串联所有处理器,传递处理结果
        return Flux.fromIterable(allProcessors)
                .reduce(Mono.just(initial), (currentResultMono, processor) ->
                        currentResultMono.flatMap(document -> {
                            log.debug("Executing processor: {}", processor.name());
                            return processor.process(document);
                        })
                )
                .flatMap(Function.identity());
    }
}

实现说明

  • 顺序执行保证:利用Flux.fromIterable遍历已排序的处理器列表,结合reduce操作符确保上一个处理器的异步任务完成后,才会执行下一个处理器
  • 结果传递逻辑:reduce的初始值是包含初始文档的Mono,每一步将当前Mono的结果传递给下一个处理器的process方法,生成新的Mono
  • 异步友好:完全基于Reactor的异步模型,不会阻塞线程,符合响应式编程的设计原则

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.02 06:20:27