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

非阻塞方式将两个Flux结果聚合为Mono的实现方案

非阻塞方式合并两个Flux并创建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;
    }
}

需要从不同方法分别获取失败账户与成功账户的Flux,非阻塞地创建Payments对象。之前尝试的代码中,调用collectList().subscribe()属于阻塞操作,且实际场景是在Rest Controller中将包含两者的对象作为响应返回,而非仅在main方法中测试。之前的尝试代码如下:

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);
    }

    public static Mono<Payments> getPaymentData() {

        Flux<Payments.SuccessAccount> accountsSucceeded = getAccountsSucceeded();
        Flux<Payments.FailedAccount> accountsFailed = getAccountsFailed();

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

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

        accountsFailed.collectList().subscribe(failedAccounts::addAll);// 这是阻塞调用

        accountsSucceeded.collectList().subscribe(successAccounts::addAll);// 这是阻塞调用

        return Mono.just(Payments.builder()
                .failedAccounts(failedAccounts)
                .successAccounts(successAccounts)
                .build());
    }
    
    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());
    }
}

非阻塞解决方案

使用Reactor的Mono.zip操作符可以并行等待多个Mono完成,全程非阻塞,完全符合响应式编程模型。

修改后的getPaymentData方法如下:

public static Mono<Payments> getPaymentData() {
    // 将Flux转换为Mono<List>,非阻塞操作
    Mono<List<Payments.SuccessAccount>> successListMono = getAccountsSucceeded().collectList();
    Mono<List<Payments.FailedAccount>> failedListMono = getAccountsFailed().collectList();

    // 并行等待两个Mono完成,合并结果构建Payments对象
    return Mono.zip(successListMono, failedListMono)
            .map(tuple -> Payments.builder()
                    .successAccounts(tuple.getT1())
                    .failedAccounts(tuple.getT2())
                    .build());
}

关键说明

  • collectList():将Flux序列转换为Mono<List>,仅在订阅触发时执行收集,全程非阻塞。
  • Mono.zip:同时订阅传入的多个Mono,等待所有Mono都完成后,将结果封装为Tuple2返回,保证并行执行且无阻塞。
  • map操作符:从Tuple2中取出两个列表,构建最终的Payments对象返回。

这种方式可以直接在Rest Controller中返回Mono<Payments>,Spring Web会自动处理订阅和响应的序列化返回,完全适配响应式Web场景。

内容的提问来源于stack exchange,提问作者nanpakal

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.08 18:45:43