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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.03 21:48:04