Spring WebFlux异步构建递归树形ChannelResponse响应的实现方案
方案解答
纯异步响应式实现方案
完全可以用纯响应式的异步逻辑实现树形结构组装,无需调用block(),核心改造思路是将递归方法的返回值改为Mono<ChannelResponse>,通过响应式操作符串联异步查询流程:
public Flux<ChannelResponse> getAll() { // 顶层节点查询后异步递归组装子树 return channelRepository.findAllByParentChannelIdOrderByOrdinality(null) .flatMap(this::buildChannelTree); } // 异步递归组装单个节点的完整树形结构 private Mono<ChannelResponse> buildChannelTree(Channel channel) { // 异步查询当前节点的子节点 return channelRepository.findAllByParentChannelIdOrderByOrdinality(channel.getId()) // 每个子节点递归组装子树 .flatMap(this::buildChannelTree) // 收集所有子节点组装结果 .collectList() // 组装当前节点的响应对象 .map(children -> new ChannelResponse(channel.getId(), children)); }
实现说明
- 全程无阻塞调用,所有数据库查询的异步结果通过
flatMap、collectList等操作符串联,完全符合WebFlux非阻塞运行规范 - 最终返回的
Flux<ChannelResponse>就是所有顶层频道的完整树形结构,可直接作为接口返回值
优化方案(避免N+1查询)
上述递归实现存在N+1查询问题,若频道数据量不大,可一次性查询全量频道数据后在内存中组装树,性能更高:
public Flux<ChannelResponse> getAll() { return channelRepository.findAll() .collectList() .flatMapMany(allChannels -> { // 构建父ID->子节点列表映射、ID->节点对象映射 Map<Long, List<ChannelResponse>> parentChildMap = new HashMap<>(); Map<Long, ChannelResponse> idNodeMap = new HashMap<>(); // 初始化所有节点的基础对象 for (Channel channel : allChannels) { ChannelResponse node = new ChannelResponse(channel.getId(), new ArrayList<>()); idNodeMap.put(channel.getId(), node); parentChildMap.computeIfAbsent(channel.getParentChannelId(), k -> new ArrayList<>()) .add(node); } // 为每个节点填充子节点列表 for (Channel channel : allChannels) { ChannelResponse node = idNodeMap.get(channel.getId()); List<ChannelResponse> children = parentChildMap.getOrDefault(channel.getId(), Collections.emptyList()); node.children().addAll(children); } // 返回所有顶层节点 return Flux.fromIterable(parentChildMap.getOrDefault(null, Collections.emptyList())); }); }
同步异步混合方案(不推荐)
如果一定要沿用原有同步递归逻辑,可将阻塞调用放到专门的阻塞调度器上执行,避免占用WebFlux的事件循环线程:
public Flux<ChannelResponse> getAll() { return channelRepository.findAllByParentChannelIdOrderByOrdinality(null) .flatMap(channel -> Mono.fromCallable(() -> getChannelDataRecursive(channel)) // 放到专门处理阻塞任务的boundedElastic调度器执行 .subscribeOn(Schedulers.boundedElastic())); }
注意事项
该方案本质还是阻塞调用,仅适合临时过渡,长期来看还是建议使用纯响应式实现。
内容的提问来源于stack exchange,提问作者Andrew Lalis
相关产品推荐
相关产品推荐

