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

基于Spring WebFlux构建非阻塞Dashboard REST API的问题求助

构建非阻塞Dashboard REST API的响应式实现问题

需求说明

  • 构建非阻塞的Dashboard REST API,数据来自三个独立源:
    1. 第三方API提供的Customer Profile
    2. 另一第三方API提供的Customer Accounts
    3. 数据库提供的Customer Beneficiary List
  • 响应结构需包含profile、accounts、beneficiaries三个字段
  • 核心要求:任一数据源返回结果后,立即返回更新后的完整响应对象,无需等待其他数据源完成

现有尝试代码

public Flux<CustomerDashboardDTO> getCustomerDashboard(String customerId) {

    // Create a Mono for profile and beneficiaries
    Mono<CustomerProfile> profileMono = Mono.just(createDummyProfile());
    Mono<List<BeneficiaryDetails>> beneficiariesMono = Mono.just(createBene());
    // Create a Flux to emit the initial data
    return Mono.zip(profileMono, beneficiariesMono)
          .flatMapMany(tuple -> {
              // Create the dashboardDTO and set profile and beneficiaries
              CustomerDashboardDTO dashboardDTO = new CustomerDashboardDTO();
              dashboardDTO.setProfile(tuple.getT1());
              dashboardDTO.setBeneficiaries(tuple.getT2());
              

              // Emit the initial DTO immediately
              Flux<CustomerDashboardDTO> initialEmission = Flux.just(dashboardDTO);

              // After a delay, fetch accounts and update the DTO
              Mono<CustomerDashboardDTO> updatedEmission = Mono.delay(Duration.ofSeconds(10))
                    .flatMap(delay -> {
                        CustomerAccount account = createDummyAccounts(); // Fetch accounts
                        dashboardDTO.setAccounts(account); // Update DTO with accounts
                        return Mono.just(dashboardDTO); // Emit the updated DTO
                    });

              // Combine initial and updated emissions
              return initialEmission.concatWith(updatedEmission);
          })
          .doOnNext(dto -> System.out.println("Emitting DTO: " + dto)); // Log emission
}

当前代码问题

  1. 初始响应必须等待profile和beneficiaries两个数据源都完成才会发射,不符合“数据到即返回”的要求
  2. accounts数据被硬编码延迟10秒处理,且依赖初始DTO创建后的流程,无法独立触发更新

解决方案

要实现“任一数据源返回即更新响应”的核心需求,需要让三个数据源的处理完全独立,每个数据源完成时立即更新并发射最新的DTO。以下是修正后的实现:

关键思路

  • 为每个数据源创建独立的Mono,各自异步获取数据
  • 使用原子引用保存当前最新的DTO状态,避免并发修改问题
  • 每个数据源完成时,基于当前状态生成新的DTO实例并更新状态,然后发射该DTO
  • 合并三个数据源的更新流,确保任一数据源完成就立即发射更新后的响应

修正后的代码

import java.util.concurrent.atomic.AtomicReference;
import reactor.core.publisher.Flux;
import reactor.core.publisher.Mono;

public Flux<CustomerDashboardDTO> getCustomerDashboard(String customerId) {
    // 1. 定义三个独立的数据源Mono(替换为实际的API/DB调用)
    Mono<CustomerProfile> profileMono = fetchProfileFromThirdParty(customerId);
    Mono<List<BeneficiaryDetails>> beneficiariesMono = fetchBeneficiariesFromDB(customerId);
    Mono<CustomerAccount> accountsMono = fetchAccountsFromThirdParty(customerId);

    // 2. 初始化原子引用保存当前DTO状态,初始为空白对象
    AtomicReference<CustomerDashboardDTO> currentDtoRef = new AtomicReference<>(new CustomerDashboardDTO());

    // 3. 为每个数据源创建更新流:数据返回时生成新的DTO并发射
    Flux<CustomerDashboardDTO> profileUpdates = profileMono.map(profile -> {
        CustomerDashboardDTO updatedDto = new CustomerDashboardDTO(currentDtoRef.get());
        updatedDto.setProfile(profile);
        currentDtoRef.set(updatedDto);
        return updatedDto;
    });

    Flux<CustomerDashboardDTO> beneficiaryUpdates = beneficiariesMono.map(beneficiaries -> {
        CustomerDashboardDTO updatedDto = new CustomerDashboardDTO(currentDtoRef.get());
        updatedDto.setBeneficiaries(beneficiaries);
        currentDtoRef.set(updatedDto);
        return updatedDto;
    });

    Flux<CustomerDashboardDTO> accountUpdates = accountsMono.map(account -> {
        CustomerDashboardDTO updatedDto = new CustomerDashboardDTO(currentDtoRef.get());
        updatedDto.setAccounts(account);
        currentDtoRef.set(updatedDto);
        return updatedDto;
    });

    // 4. 合并所有更新流,可选:先发射初始空白DTO(根据需求决定是否保留)
    return Flux.concat(
            Flux.just(currentDtoRef.get()), // 初始空白响应(可选)
            Flux.merge(profileUpdates, beneficiaryUpdates, accountUpdates)
        )
        .doOnNext(dto -> System.out.println("Emitting DTO: " + dto));
}

// 以下为模拟实际数据源调用的方法(替换为真实实现)
private Mono<CustomerProfile> fetchProfileFromThirdParty(String customerId) {
    // 模拟第三方API调用延迟
    return Mono.just(createDummyProfile()).delayElement(java.time.Duration.ofMillis(200));
}

private Mono<List<BeneficiaryDetails>> fetchBeneficiariesFromDB(String customerId) {
    // 模拟DB查询延迟
    return Mono.just(createBene()).delayElement(java.time.Duration.ofMillis(500));
}

private Mono<CustomerAccount> fetchAccountsFromThirdParty(String customerId) {
    // 模拟另一第三方API调用延迟
    return Mono.just(createDummyAccounts()).delayElement(java.time.Duration.ofMillis(800));
}

// 原有的模拟数据创建方法保留
private CustomerProfile createDummyProfile() { /* 实现略 */ }
private List<BeneficiaryDetails> createBene() { /* 实现略 */ }
private CustomerAccount createDummyAccounts() { /* 实现略 */ }

代码说明

  • 独立数据源处理:三个Mono各自异步执行,互不等待
  • 状态安全更新:使用AtomicReference保证DTO状态的线程安全,每次更新都创建新的DTO实例(避免并发修改问题)
  • 即时发射响应:Flux.merge会在任一数据源完成时立即发射更新后的DTO,完全符合“数据到即返回”的要求
  • 可选初始响应:如果不需要初始空白DTO,可以直接返回Flux.merge(profileUpdates, beneficiaryUpdates, accountUpdates)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.17 12:34:56