基于Spring WebFlux构建非阻塞Dashboard REST API的问题求助
构建非阻塞Dashboard REST API的响应式实现问题
需求说明
- 构建非阻塞的Dashboard REST API,数据来自三个独立源:
- 第三方API提供的Customer Profile
- 另一第三方API提供的Customer Accounts
- 数据库提供的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 }
当前代码问题
- 初始响应必须等待
profile和beneficiaries两个数据源都完成才会发射,不符合“数据到即返回”的要求 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
相关产品推荐
相关产品推荐

