Reactor响应式编程:非阻塞地聚合两个Flux流的结果为单个Mono对象


Reactor响应式编程:非阻塞地聚合两个Flux流的结果为单个Mono对象

本文旨在详细阐述在project reactor框架中,如何优雅且非阻塞地将两个独立的flux流处理后的结果聚合为一个单一的mono对象。通过分析传统阻塞式操作的弊端,我们将重点介绍并演示mono.zipwith操作符的正确使用方法,以实现高效、响应式的并发数据聚合,从而避免在异步流程中引入阻塞点。

1. 理解响应式流中的非阻塞聚合需求

在响应式编程中,我们经常需要从多个独立的异步源获取数据,并将这些数据组合成一个统一的结果对象。例如,一个支付服务可能需要同时从不同的子系统获取成功交易列表和失败交易列表,然后将它们封装在一个Payments对象中返回。

考虑以下领域模型:

package org.example;

import lombok.Builder;
import lombok.Getter;
import lombok.ToString;

import j*a.util.List;

@Getter
@Builder
@ToString
public class Payments {
    private List<SuccessAccount> successAccounts;
    private List<FailedAccount> failedAccounts;

    @Getter
    @Builder
    @ToString
    public static class SuccessAccount {
        private String name;
        private String accountNumber;
    }

    @Getter
    @Builder
    @ToString
    public static class FailedAccount {
        private String name;
        private String accountNumber;
        private String errorCode;
    }
}

假设我们有两个方法分别返回成功账户和失败账户的Flux流:

public static Flux<Payments.SuccessAccount> getAccountsSucceeded() {
    return Flux.just(Payments.SuccessAccount.builder()
                    .accountNumber("1234345")
                    .name("Payee1")
                    .build(),
            Payments.SuccessAccount.builder()
                    .accountNumber("83673674")
                    .name("Payee2")
                    .build());
}

public static Flux<Payments.FailedAccount> getAccountsFailed() {
    return Flux.just(Payments.FailedAccount.builder()
                    .accountNumber("12234345")
                    .name("Payee3")
                    .errorCode("8938")
                    .build(),
            Payments.FailedAccount.builder()
                    .accountNumber("3342343")
                    .name("Payee4")
                    .errorCode("8938")
                    .build());
}

一个常见的误区是尝试通过订阅这些Flux流并将结果收集到可变列表中,然后构建最终对象。例如:

// 这是一个阻塞的、不推荐的做法
public static Mono<Payments> getPaymentDataBlocking() {
    Flux<Payments.SuccessAccount> accountsSucceeded = getAccountsSucceeded();
    Flux<Payments.FailedAccount> accountsFailed = getAccountsFailed();

    List<Payments.SuccessAccount> successAccounts = new ArrayList<>();
    List<Payments.FailedAccount> failedAccounts = new ArrayList<>();

    // 调用 subscribe() 会立即触发流的执行,并在当前线程等待结果,导致阻塞
    accountsFailed.collectList().subscribe(failedAccounts::addAll);
    accountsSucceeded.collectList().subscribe(successAccounts::addAll);

    return Mono.just(Payments.builder()
            .failedAccounts(failedAccounts)
            .successAccounts(successAccounts)
            .build());
}

上述代码中的subscribe()调用是阻塞的,因为它会在当前线程等待collectList()操作完成,这违背了Reactor非阻塞的原则。在实际的Web服务或异步处理场景中,这种阻塞操作会导致线程池资源耗尽,严重影响系统吞吐量和响应性。

2. 使用Mono.zipWith 实现非阻塞聚合

为了在Reactor中实现真正的非阻塞聚合,我们需要利用其提供的组合操作符。Mono.zipWith(或Mono.zip)是解决此类问题的理想选择。它允许我们将两个Mono(或多个Mono)的结果组合起来,一旦所有源Mono都完成了并产生了它们的值,就会使用一个提供的BiFunction(或Function)来处理这些值,并生成一个新的Mono结果。

Decktopus AI Decktopus AI

AI在线生成高质量演示文稿

Decktopus AI 153 查看详情 Decktopus AI

具体步骤如下:

  1. 将Flux转换为Mono 首先,我们需要将每个Flux流通过collectList()操作符转换为一个发出单个List的Mono。这个Mono将在原始Flux完成并收集所有元素后发出其列表。
  2. 使用zipWith组合: 接下来,将第一个Mono与第二个Mono使用zipWith操作符进行组合。
  3. 提供组合函数: zipWith需要一个BiFunction作为参数,该函数接收两个Mono发出的值(即两个List),并返回我们期望的最终结果(即Payments对象)。

下面是使用Mono.zipWith实现的非阻塞解决方案:

package org.example;

import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

import j*a.util.ArrayList;
import j*a.util.List;

public class Main {
    public static void main(String[] args) {
        // 订阅并打印结果,这是在应用程序入口点进行的操作,不会阻塞核心业务逻辑
        getPaymentData().subscribe(System.out::println);

        // 为了在main方法中观察异步结果,通常需要一些延迟或等待机制
        // 在实际应用中,例如Spring WebFlux控制器,Mono会被框架自动订阅和处理
        try {
            Thread.sleep(1000); // 简单等待,仅用于演示
        } catch (InterruptedException e) {
            Thread.currentThread().interrupt();
        }
    }

    public static Mono<Payments> getPaymentData() {
        Flux<Payments.SuccessAccount> accountsSucceededFlux = getAccountsSucceeded();
        Flux<Payments.FailedAccount> accountsFailedFlux = getAccountsFailed();

        // 将Flux转换为Mono<List>
        Mono<List<Payments.SuccessAccount>> successAccountsMono = accountsSucceededFlux.collectList();
        Mono<List<Payments.FailedAccount>> failedAccountsMono = accountsFailedFlux.collectList();

        // 使用 zipWith 组合两个 Mono 的结果
        Mono<Payments> combinedPaymentsMono = failedAccountsMono.zipWith(
                successAccountsMono,
                (failedAccounts, successAccounts) -> Payments.builder()
                        .failedAccounts(failedAccounts)
                        .successAccounts(successAccounts)
                        .build()
        );

        return combinedPaymentsMono;
    }

    public static Flux<Payments.SuccessAccount> getAccountsSucceeded() {
        return Flux.just(Payments.SuccessAccount.builder()
                        .accountNumber("1234345")
                        .name("Payee1")
                        .build(),
                Payments.SuccessAccount.builder()
                        .accountNumber("83673674")
                        .name("Payee2")
                        .build());
    }

    public static Flux<Payments.FailedAccount> getAccountsFailed() {
        return Flux.just(Payments.FailedAccount.builder()
                        .accountNumber("12234345")
                        .name("Payee3")
                        .errorCode("8938")
                        .build(),
                Payments.FailedAccount.builder()
                        .accountNumber("3342343")
                        .name("Payee4")
                        .errorCode("8938")
                        .build());
    }
}

在这个改进后的getPaymentData()方法中:

  • accountsSucceededFlux.collectList()和accountsFailedFlux.collectList()各自返回一个Mono。这两个Mono会并行地收集它们各自Flux中的所有元素。
  • failedAccountsMono.zipWith(successAccountsMono, ...)操作符会等待这两个Mono都完成并发出它们的结果(即两个List)。
  • 一旦两个List都可用,zipWith会调用提供的BiFunction,将这两个List作为参数传入,然后使用它们来构建并发出最终的Payments对象。
  • 整个过程都是非阻塞的,getPaymentData()方法会立即返回一个Mono,而实际的数据处理和对象构建则会在背后的Reactor调度器上异步执行。

3. 注意事项与最佳实践

  • 避免中间订阅: 在响应式链中,除了最终的消费者(如REST控制器返回Mono或在main方法中打印结果),应尽量避免使用subscribe()来获取中间结果。subscribe()会触发流的执行,并且其副作用(如修改外部变量)在异步环境中难以管理,也容易引入阻塞。
  • 利用组合操作符: Reactor提供了丰富的组合操作符(如zip、merge、concat、when等),它们是处理多个响应式流的强大工具。选择正确的操作符取决于你希望如何组合这些流的行为(例如,并行等待所有完成、按顺序合并、或只关心第一个完成的)。
  • 错误处理: zipWith操作符具有短路特性。如果其中任何一个源Mono发出错误,那么zipWith返回的Mono也会立即发出相同的错误,而不会等待其他源完成。这对于快速失败和错误传播非常有用。
  • 可读性和可维护性: 保持响应式链的流畅性,避免将异步操作拆分为多个独立的阻塞步骤,可以显著提高代码的可读性和可维护性。

总结

通过Mono.zipWith操作符,我们能够优雅且高效地在Project Reactor中聚合来自多个Flux流的异步结果,并将其封装成一个单一的Mono对象。这种模式是构建高性能、非阻塞响应式应用程序的关键,它确保了在处理并发数据源时,应用程序能够充分利用资源并保持出色的响应能力。理解并正确运用这些组合操作符,是掌握Reactor响应式编程范式的核心。

以上就是Reactor响应式编程:非阻塞地聚合两个Flux流的结果为单个Mono对象的详细内容,更多请关注其它相关文章!


# 这是  # 网站系统建设招标公告  # 抖音视频营销推广意图吗  # 广州市口碑seo费用  # 即墨区网站优化培训  # 重庆安徽网站建设  # 深泽品牌网站推广技巧  # 怎样做个人微信营销推广  # 房产网站建设厂家  # 东营seo优化优势  # 夫唯seo 124期  # 就会  # 加载  # react  # 如何实现  # 并将  # 应用程序  # 第一个  # 转换为  # 这两个  # 多个  # 响应式编程  # ai  # 工具  # java 


相关栏目: 【 Google疑问12 】 【 Facebook疑问10 】 【 优化推广96088 】 【 技术知识133117 】 【 IDC资讯59369 】 【 网络运营7196 】 【 IT资讯61894


相关推荐: 微信客户端如何找回密码_微信客户端忘记密码找回方法  海棠阅读网页版_进入海棠网页版在线阅读中心  Win10显卡驱动安装失败怎么办 Win10使用DDU彻底卸载驱动【解决】  解决异步Python机器人中同步操作的阻塞问题  mail.qq.com登录入口 QQ邮箱网页版直达  《战地6》反作弊已成功拦截240万次作弊 发售第一周98%比赛没有作弊  yandex网页版直接登录 yandex官方入口平台访问方法  mysql中如何分析索引使用情况_mysql索引使用分析方法  Sublime怎么自动添加CSS前缀_Sublime安装Autoprefixer插件  冬季去寒冷地区旅游,以下哪种做法有助于缓解冻伤  如何查询个人病历记录  附近酒吧怎么找?  视频号视频怎么免费保存到相册?保存到相册需要注意什么?  vivo手机视频通话美颜怎么设置_vivo视频通话美颜开启方法  Excel怎么用XLOOKUP函数实现双向查找_ExcelXLOOKUP替代VLOOKUP+HLOOKUP的高级用法  J*a中逻辑运算符如何使用_逻辑与或非的基础用法讲解  PPT页面尺寸怎么修改 PPT自定义幻灯片大小与方向设置【教程】  PHP动态导航按钮:根据用户登录状态切换链接与文本  《友玩*》创建群聊方法  Lar*el如何创建自定义的辅助函数(Helpers)_Lar*el全局函数定义与加载方法  顺丰快递收费标准查询_如何查看顺丰最新收费价格  《大周列国志》皇帝律令功能介绍  解决 Vue 3 组件未定义错误:理解 createApp 与根组件的正确使用  顺丰快递怎么查物流_顺丰快递物流信息实时查询操作指南  解决CSS容器溢出问题:使用calc()实现精确布局与边距控制  虫虫漫画排行榜单入口_虫虫漫画编辑推荐入口  画质怪兽120帧安卓和平精英免费版  在Peewee中处理PostgreSQL记录重复:一站式数据摄取教程  POKI小游戏在线免费入口链接 POKI小游戏无下载秒玩玩  作业帮网页版不用下载入口 在线问老师快速答疑  创建快捷方式启动系统保护  Linux如何开发轻量级数据服务模块_Linux服务化设计  漫蛙manwa官网浏览入口_漫蛙漫画网页版访问链接  狙击外星人小游戏在线链接_狙击外星人小游戏网页链接  铁拳8在线玩 铁拳8在线秒玩入口  b站怎么设置动态仅粉丝可见_b站动态粉丝可见设置方法  CSS动画如何实现图标旋转并放大_transform rotate scale @keyframes实现  谷歌邮箱官方入口链接 谷歌邮箱网页版电脑端快速登录  PHP utf8_encode 字符编码转换疑难解析与最佳实践  Sublime怎么快速复制文件路径_Sublime右键菜单增强技巧  ExcelSCAN与LAMBDA如何创建自定义移动平均函数_SCAN实现任意窗口期移动平均计算  解决SQLAlchemy模型跨文件关联的Linter兼容性指南  食品生产用水只要符合国家规定的生活饮用水卫生标准就可以吗  怎样让Windows 11的开始菜单恢复经典样式_Open-Shell工具使用指南【怀旧】  PHP utf8_encode 字符编码转换陷阱与解决方案  韩小圈网页版PC端入口 韩小圈网页版官方网站入口  Sublime怎么格式化HTML代码_Sublime前端代码美化插件使用指南  yy漫画官方网站登录入口_yy漫画在线阅读页面地址  J*aScript 数值去小数位处理:多种方法与实践  《兴业银行》注册登录方法 

 2025-12-03

了解您产品搜索量及市场趋势,制定营销计划

同行竞争及网站分析保障您的广告效果

点击免费数据支持

提交您的需求,1小时内享受我们的专业解答。

运城市盐湖区信雨科技有限公司


运城市盐湖区信雨科技有限公司

运城市盐湖区信雨科技有限公司是一家深耕海外推广领域十年的专业服务商,作为谷歌推广与Facebook广告全球合作伙伴,聚焦外贸企业出海痛点,以数字化营销为核心,提供一站式海外营销解决方案。公司凭借十年行业沉淀与平台官方资源加持,打破传统外贸获客壁垒,助力企业高效开拓全球市场,成为中小企业出海的可靠合作伙伴。

 8156699

 13765294890

 8156699@qq.com

Notice

We and selected third parties use cookies or similar technologies for technical purposes and, with your consent, for other purposes as specified in the cookie policy.
You can consent to the use of such technologies by closing this notice, by interacting with any link or button outside of this notice or by continuing to browse otherwise.