如何使用Reactor实现顺序执行多Consumer遇false即终止的逻辑
Reactor 消费链顺序执行并提前终止实现方案
实现思路
你的命令式逻辑核心是顺序执行消费器、布尔返回值控制是否终止流程,Reactor 可通过 concatMap + takeWhile 组合完美对齐该逻辑:
concatMap保证消费器严格按列表顺序执行,不会并发乱序takeWhile遇到返回值为false时直接终止整个流,符合提前终止的需求
代码实现
首先定义和你原有逻辑对齐的消费接口:
interface MyConsumer { // 入参为待处理的值,返回值标识是否继续执行后续消费器 boolean consume(Object value); }
核心响应式实现代码:
import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import java.util.List; public class ConsumerChainDemo { public static void main(String[] args) { // 待处理的原始值 Object value = 1; // 消费器列表,和你原有逻辑的 consumers 完全对齐 List<MyConsumer> consumers = List.of( v -> { System.out.println("消费器1执行,值为:" + v); return true; }, v -> { System.out.println("消费器2执行,值为:" + v); return true; }, v -> { System.out.println("消费器3执行,值为:" + v); return false; // 返回false,后续消费器不会执行 }, v -> { System.out.println("消费器4执行,值为:" + v); return true; } ); // 核心逻辑 Flux.fromIterable(consumers) // 顺序执行每个消费器的处理逻辑,同步异步都支持 .concatMap(consumer -> Mono.fromCallable(() -> consumer.consume(value))) // 返回true继续执行,遇到false立刻终止流 .takeWhile(continueFlag -> continueFlag) // 触发流执行(Web 场景下直接返回流给框架即可,无需调用block) .blockLast(); } }
适配说明
如果你的消费逻辑本身是异步IO操作,只需将 MyConsumer 接口的返回值改为 Mono<Boolean>,核心逻辑不需要调整,直接替换 concatMap 内的代码为 consumer.consume(value) 即可。
内容的提问来源于stack exchange,提问作者user11258188
相关产品推荐
相关产品推荐

