Spring Webflux非阻塞实现:Mono集合转HashMap方案问询
问题描述
我是Spring Webflux的新手,正在编写一个聚合服务,调用两个不支持批量查询的遗留接口:
- 第一个接口返回用户的所有账户(
Mono响应),如果调用失败,流程终止,无需查询第二个接口; - 第二个接口仅接受单个账户编号,并返回该账户的交易记录(
Mono响应); - 交易记录的响应不包含账户标识,目标是生成如下结构的响应:
@Value(staticConstructor = "of") public class MyResponse { Map<String, List<Transaction>> accountTransactions; }
当前服务实现如下:
@RequiredArgsConstructor public class MyService { private final UserClient userClient; private final TransactionsClient transactionsClient; public Mono<MyResponse> getTransactions(..some relavant params..) { return userClient.getAccounts(..some relavant params..) .zipWhen(accounts -> getTransactions(account)) .map(Tuple2::getT2) .subscribeOn(boundedElastic()); } public Mono<MyResponse> getTransactions(Accounts accounts) { Map<String, List<Transaction>> transactionsOfAccount = new HashMap<>(); for (var account : accounts) { transactionsOfAccount.put(account.getName(), transactionsClient.getTransactions(account) .map(r -> {..some transformation..}) .subscribeOn(boundedElastic()) .block()); // 不希望用block() } return Mono.just(MyResponse.of(transactionsOfAccount)); // 不希望这样封装 } }
请问有没有非阻塞的实现方式?
提前感谢!
非阻塞实现方案
完全可以通过Webflux的响应式操作符实现非阻塞逻辑,核心是用响应式流的批量处理替代同步循环+block():
@RequiredArgsConstructor public class MyService { private final UserClient userClient; private final TransactionsClient transactionsClient; public Mono<MyResponse> getTransactions(..some relevant params..) { return userClient.getAccounts(..some relevant params..) // 将账户集合转为Flux,实现逐个处理的响应式流 .flatMapMany(accounts -> Flux.fromIterable(accounts)) // 为每个账户异步查询交易记录,同时保留账户与交易的关联 .flatMap(account -> transactionsClient.getTransactions(account) .map(r -> {..some transformation..}) .subscribeOn(boundedElastic()) // 用Entry配对账户名和交易记录,解决交易无账户标识的问题 .map(transactions -> new AbstractMap.SimpleEntry<>(account.getName(), transactions)) ) // 将所有配对结果收集为目标Map,响应式聚合无阻塞 .collectMap( AbstractMap.SimpleEntry::getKey, AbstractMap.SimpleEntry::getValue ) // 封装为最终响应对象 .map(MyResponse::of) .subscribeOn(boundedElastic()); } }
关键逻辑说明
flatMapMany:把Mono<Accounts>转换为Flux<Account>,将批量账户转为逐个处理的响应式流flatMap:异步调用交易接口,同时通过SimpleEntry绑定账户名和交易记录,确保后续能正确关联collectMap:响应式的聚合操作,自动收集所有结果为目标Map,全程无阻塞- 错误处理:如果第一个账户接口调用失败,整个流会直接终止,不会触发后续的交易接口查询,符合需求
另外原代码中的zipWhen属于冗余操作,直接在账户流的后续链路中完成交易查询即可。
内容的提问来源于stack exchange,提问作者asli
相关产品推荐
相关产品推荐

