Reactor Java中如何用非阻塞Mono实现用户任务关联账户更新逻辑
正确实现代码
public Mono<Void> execute(List<String> users) { return Flux.fromIterable(users) // 异步查询用户关联的任务列表 .flatMap(jobDao::getJobsByUser) // 把任务列表打平为单个JobInfo的流 .flatMapIterable(jobInfos -> jobInfos) // 异步查询任务关联的账户列表 .flatMap(jobInfo -> accountDam.getAccounts(jobInfo.getAccount())) // 把账户列表打平为单个Account的流 .flatMapIterable(accounts -> accounts) // 执行账户更新逻辑 .flatMap(account -> { // 若更新为异步方法,直接返回对应Mono<Void>即可 // 若更新为同步逻辑,用Mono.fromRunnable包裹: // return Mono.fromRunnable(() -> { /* 你的更新代码 */ }); return updateAccount(account); }) // 所有操作完成后返回结束信号 .then(); }
核心改造说明
原代码报错的核心原因是Reactor的reactor-http-nio线程属于非阻塞线程,禁止调用block()这类主动阻塞等待结果的方法,违反了响应式编程的非阻塞规范。改造后的逻辑全程采用异步链式调用:
- 用
Flux.fromIterable把普通用户集合转为响应式流,适配非阻塞操作 flatMap用于处理返回Mono的异步调用,不会阻塞等待结果,结果就绪后自动触发后续流程flatMapIterable直接把异步返回的集合打平为单个元素的流,不需要手动遍历集合- 若需要控制并发避免下游压力过大,可以给
flatMap增加第二个参数指定最大并发数,例如flatMap(jobDao::getJobsByUser, 10)表示最多同时执行10个任务查询请求
内容的提问来源于stack exchange,提问作者Delantha
相关产品推荐
相关产品推荐

