0

0

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

DDD

DDD

发布时间:2025-12-03 16:57:12

|

905人浏览过

|

来源于php中文网

原创

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 successAccounts;
    private List 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 getAccountsSucceeded() {
    return Flux.just(Payments.SuccessAccount.builder()
                    .accountNumber("1234345")
                    .name("Payee1")
                    .build(),
            Payments.SuccessAccount.builder()
                    .accountNumber("83673674")
                    .name("Payee2")
                    .build());
}

public static Flux 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 getPaymentDataBlocking() {
    Flux accountsSucceeded = getAccountsSucceeded();
    Flux accountsFailed = getAccountsFailed();

    List successAccounts = new ArrayList<>();
    List 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结果。

析易-AI论文_数据分析
析易-AI论文_数据分析

一个专业的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 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 getPaymentData() {
        Flux accountsSucceededFlux = getAccountsSucceeded();
        Flux accountsFailedFlux = getAccountsFailed();

        // 将Flux转换为Mono
        Mono> successAccountsMono = accountsSucceededFlux.collectList();
        Mono> failedAccountsMono = accountsFailedFlux.collectList();

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

        return combinedPaymentsMono;
    }

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

    public static Flux 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响应式编程范式的核心。

相关文章

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

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

下载

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

相关专题

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

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

481

2023.08.10

function是什么
function是什么

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

478

2023.08.04

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

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

163

2023.10.07

高德地图升级方法汇总
高德地图升级方法汇总

本专题整合了高德地图升级相关教程,阅读专题下面的文章了解更多详细内容。

68

2026.01.16

全民K歌得高分教程大全
全民K歌得高分教程大全

本专题整合了全民K歌得高分技巧汇总,阅读专题下面的文章了解更多详细内容。

123

2026.01.16

C++ 单元测试与代码质量保障
C++ 单元测试与代码质量保障

本专题系统讲解 C++ 在单元测试与代码质量保障方面的实战方法,包括测试驱动开发理念、Google Test/Google Mock 的使用、测试用例设计、边界条件验证、持续集成中的自动化测试流程,以及常见代码质量问题的发现与修复。通过工程化示例,帮助开发者建立 可测试、可维护、高质量的 C++ 项目体系。

34

2026.01.16

java数据库连接教程大全
java数据库连接教程大全

本专题整合了java数据库连接相关教程,阅读专题下面的文章了解更多详细内容。

39

2026.01.15

Java音频处理教程汇总
Java音频处理教程汇总

本专题整合了java音频处理教程大全,阅读专题下面的文章了解更多详细内容。

19

2026.01.15

windows查看wifi密码教程大全
windows查看wifi密码教程大全

本专题整合了windows查看wifi密码教程大全,阅读专题下面的文章了解更多详细内容。

85

2026.01.15

热门下载

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

精品课程

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

共58课时 | 3.8万人学习

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

共12课时 | 1.0万人学习

React核心原理新老生命周期精讲
React核心原理新老生命周期精讲

共12课时 | 1万人学习

关于我们 免责申明 举报中心 意见反馈 讲师合作 广告合作 最新更新
php中文网:公益在线php培训,帮助PHP学习者快速成长!
关注服务号 技术交流群
PHP中文网订阅号
每天精选资源文章推送

Copyright 2014-2026 https://www.php.cn/ All Rights Reserved | php.cn | 湘ICP备2023035733号