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

Java 11+Project Reactor中Mono的flatMap引发方法重复调用问题

问题:使用flatMap后deleteAccount方法被调用两次

我正在实现一个返回Mono<Response>的deleteAccount方法,该方法依赖另一个服务的accountDataProvider.accountDataFor(accountId)(同样返回Mono<Response>)。为了避免在管道末尾使用block()实现异步处理,我把map替换成flatMap后,发现deleteAccount方法被调用了两次。


相关代码

原使用map的实现(含block())

public Mono<Response> deleteAccount(String accountId) {
    return accountDataProvider.accountDataFor(accountId)
            .map(accountData -> {
                if (accountData.isActive) {
                    return deleteAccount(accountId);
                }
                return Mono.just(new Response());
            }).block();
}

改用flatMap的实现

public Mono<Response> deleteAccount(String accountId) {
    return accountDataProvider.accountDataFor(accountId)
            .flatMap(accountData -> {
                if (accountData.isActive) {
                    return deleteAccount(accountId);
                }
                return Mono.just(new Response());
            });
}

deleteAccount内部逻辑

var response = new Response();
accountRegistry.getSessions(List.of(accountId)).stream()
            .iterator()
            .forEachRemaining(session -> {
                    response.addToDeletedSessions(session);
            });
return Mono.just(response);

测试情况

测试代码中执行var res = sut.deleteAccount(accountId);后两次调用block(),验证时发现accountRegistry.getSessions被执行了2次,而预期是1次,报错:

期望调用1次:
实际调用2次:


问题分析与解决

核心问题

  1. 递归调用错误:改用flatMap的实现中,当accountData.isActive为true时,直接递归调用了deleteAccount(accountId),这会导致每次订阅都会触发完整的管道流程,包括再次调用accountDataProvider.accountDataFor和删除逻辑。
  2. Mono的冷发布特性:Reactor中的Mono是冷发布者,每次调用block()(即订阅)都会重新执行整个管道内的逻辑。测试中两次调用block(),相当于触发了两次完整执行。

解决方案

1. 拆分递归逻辑,避免自我调用

把实际的删除逻辑抽成独立方法,不再递归调用整个deleteAccount:

// 独立的实际删除逻辑
private Mono<Response> doDeleteAccount(String accountId) {
    var response = new Response();
    accountRegistry.getSessions(List.of(accountId)).stream()
                .iterator()
                .forEachRemaining(session -> {
                        response.addToDeletedSessions(session);
                });
    return Mono.just(response);
}

// 修正后的deleteAccount方法
public Mono<Response> deleteAccount(String accountId) {
    return accountDataProvider.accountDataFor(accountId)
            .flatMap(accountData -> {
                if (accountData.isActive) {
                    return doDeleteAccount(accountId); // 调用独立删除方法
                }
                return Mono.just(new Response());
            });
}

2. 测试中避免多次订阅

如果测试需要多次获取结果,使用cache()操作符缓存结果,确保多次订阅只执行一次逻辑:

var res = sut.deleteAccount(accountId).cache();
res.block(); // 第一次执行完整逻辑并缓存结果
res.block(); // 第二次直接使用缓存,不重复执行

或者直接一次性获取结果后多次验证:

var response = sut.deleteAccount(accountId).block();
// 基于已获取的response做多次验证,无需再次调用block()

额外说明

原map实现本身存在逻辑问题:map返回Mono<Response>会让管道变成Mono<Mono<Response>>,后续的block()只是取出内部的Mono,整个方法是同步阻塞的,完全失去了Reactor异步的意义。改用flatMap是正确的方向,但递归调用是导致重复执行的关键诱因。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.12 08:52:28