如何用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
相关产品推荐
相关产品推荐

