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

大丽小哥_7254

大丽小哥_7254

2025-12-03

905人浏览

原创

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

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

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

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

考虑以下领域模型:

package org.example;

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

import java.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;
    }
}</failedaccount></successaccount>

假设我们有两个方法分别返回成功账户和失败账户的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());
}</payments.failedaccount></payments.successaccount>

一个常见的误区是尝试通过订阅这些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());
}</payments.failedaccount></payments.successaccount></payments.failedaccount></payments.successaccount></payments>

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

Alibabacloud Sdk Client Initialization For Java
Alibabacloud Sdk Client Initialization For Java

在 Java 中初始化和管理阿里云 SDK客户端。包括单例模式、线程安全、endpoint 与 region 配置、VPC 终端节点、同步与异步等。

下载

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

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

具体步骤如下:

  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 java.util.ArrayList;
import java.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>> successAccountsMono = accountsSucceededFlux.collectList();
        Mono<list>> 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());
    }
}</payments.failedaccount></payments.successaccount></payments></list></list></list></payments.failedaccount></payments.successaccount></payments>

在这个改进后的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响应式编程范式的核心。

相关文章

编程速学教程(入门课程)
编程速学教程(入门课程)

编程怎么学习?编程怎么入门?编程在哪学?编程怎么学才快?不用担心,这里为大家提供了编程速学教程(入门课程),有需要的小伙伴保存下载就能学习啦!

下载

相关标签:

react java 响应式编程

本站声明:本文内容由网友自发贡献,版权归原作者所有,本站不承担相应法律责任。如您发现有涉嫌抄袭侵权的内容,请联系admin@php.cn

相关专题

更多
线程和进程的区别
线程和进程的区别

线程和进程的区别:线程是进程的一部分,用于实现并发和并行操作,而线程共享进程的资源,通信更方便快捷,切换开销较小。本专题为大家提供线程和进程区别相关的各种文章、以及下载和课程。

2023.08.10

3398

6

function是什么
function是什么

function是函数的意思,是一段具有特定功能的可重复使用的代码块,是程序的基本组成单元之一,可以接受输入参数,执行特定的操作,并返回结果。本专题为大家提供function是什么的相关的文章、下载、课程内容,供大家免费下载体验。

2023.08.04

2440

5

js函数function用法
js函数function用法

js函数function用法有:1、声明函数;2、调用函数;3、函数参数;4、函数返回值;5、匿名函数;6、函数作为参数;7、函数作用域;8、递归函数。本专题提供js函数function用法的相关文章内容,大家可以免费阅读。

2023.10.07

414

5

Aionclaw智能助手介绍
Aionclaw智能助手介绍

本专题汇总了AionClaw(AI龙虾助手)的功能介绍与在线使用入口。AionClaw是杭州趣猿人工智能有限公司推出的桌面级AI智能体,能直接在电脑上读写文件、运行脚本、操作浏览器,自动交付Word、PPT、Excel等成品。

2026.09.20

20

13

AionClaw AI智能体与电脑自动化任务执行功能使用教程
AionClaw AI智能体与电脑自动化任务执行功能使用教程

AionClaw专题整理AI智能体与电脑自动化相关功能使用教程,涵盖安装部署、AI任务执行、Skills技能、文件处理、浏览器控制、电脑操作、持久记忆、聊天工具连接以及办公、编程和内容创作等功能,帮助用户快速掌握AionClaw的实际使用方法。

2026.09.20

0

15

AI视频生成软件推荐
AI视频生成软件推荐

本专题汇总了当前主流的AI视频生成软件推荐与排行榜单,涵盖seko、AniShort、剧云、Lovart、LiblibAI及立刻mv等热门工具。同时整理了各软件在文生视频、图生视频、时长限制、画质表现及免费额度等方面的差异对比,助您快速选对适合创作需求的AI视频生成工具。

2026.09.16

180

9

ai生成视频的工具免费版合集
ai生成视频的工具免费版合集

本专题汇总了当前免费AI生成视频工具的排行榜与推荐清单,涵盖seko、讯飞智作、AniShort及剧云、Lovart等多模型集成平台。同时整理了各工具的免费额度、输出时长、水印政策及适用场景差异,助您快速选择合适工具开启AI视频创作。

2026.09.16

80

10

Pandas时间序列分析与可视化报表
Pandas时间序列分析与可视化报表

本专题整理Pandas日期转换、时间索引、重采样、滚动窗口、时区处理、plot绘图、Styler表格样式和报表输出方法。

2026.09.16

80

23

Pandas数据筛选索引与清洗处理
Pandas数据筛选索引与清洗处理

本专题整理Pandas中的loc、iloc、条件筛选、query查询、缺失值处理、重复值删除、类型转换和字符串列清洗方法。

2026.09.16

60

25

热门下载

更多
网站特效
/
网站源码
/
网站素材
/
前端模板

精品课程

更多
相关推荐
/
热门推荐
/
最新课程
React 教程
React 教程

共58课时 | 11.9万人学习

国外Web开发全栈课程全集
国外Web开发全栈课程全集

共12课时 | 1.4万人学习