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次:
问题分析与解决
核心问题
- 递归调用错误:改用flatMap的实现中,当
accountData.isActive为true时,直接递归调用了deleteAccount(accountId),这会导致每次订阅都会触发完整的管道流程,包括再次调用accountDataProvider.accountDataFor和删除逻辑。 - 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
相关产品推荐
相关产品推荐

