基于Project Reactor,使用DFS判断子节点处理是否全部完成
嘿,作为Reactor新手碰到这个问题真的太常见啦!我来帮你搞定怎么准确捕捉到父节点所有子节点处理完成的时机~
首先得明确你遇到的核心问题:expandDeep的工作机制是拉取式发射元素——它会先把父节点发射出来,紧接着就认为父节点的流“完成”了,然后再去拉取并处理子节点。这就是为什么父节点的doOnComplete会在子节点处理完前触发的原因,它和你期望的“父节点等所有子节点处理完再标记完成”的逻辑不匹配。
先把你描述的树结构梳理清楚,方便后续说明:
- A
- AA
- AAA
- AAB
- AAC
- AB
- AC
- AD
- AE
- AEA
- AEB
- AA
解决方案:递归串联处理(推荐)
这种方式完全贴合DFS的逻辑,能确保父节点的完成回调严格在所有子节点处理完成后触发。核心思路是:先处理当前节点,再递归处理所有子节点,只有当所有子节点的处理流都完成后,当前父节点的流才会标记完成。
1. 补全节点结构代码
先把你的Node类补全(包含模拟处理逻辑):
import lombok.Builder; import lombok.Data; import lombok.Singular; import reactor.core.publisher.Flux; import reactor.core.publisher.Mono; import java.util.List; @Data @Builder public class Node { private String name; @Singular private List<Node> children; // 模拟每个节点的耗时处理逻辑,方便观察顺序 public Mono<Void> process() { return Mono.fromRunnable(() -> System.out.println("Processing: " + name)) .delayElement(java.time.Duration.ofMillis(100)); } }
2. DFS处理实现
public class TreeDFSProcessor { // 递归处理树,严格保证DFS顺序,父节点等所有子节点处理完才完成 public static Flux<Void> processTreeDFS(Node root) { // 先处理当前节点,再处理所有子节点,最后触发父节点的完成回调 return root.process() .thenMany(Flux.fromIterable(root.getChildren()) // 并发数设为1,确保深度优先顺序(否则会并行处理子节点) .flatMap(TreeDFSProcessor::processTreeDFS, 1)) .doOnComplete(() -> System.out.println("✅ All children of " + root.getName() + " are fully processed!")); } public static void main(String[] args) { // 构建你描述的树结构 Node tree = Node.builder() .name("A") .child(Node.builder() .name("AA") .child(Node.builder().name("AAA").build()) .child(Node.builder().name("AAB").build()) .child(Node.builder().name("AAC").build()) .build()) .child(Node.builder().name("AB").build()) .child(Node.builder().name("AC").build()) .child(Node.builder().name("AD").build()) .child(Node.builder() .name("AE") .child(Node.builder().name("AEA").build()) .child(Node.builder().name("AEB").build()) .build()) .build(); // 阻塞等待处理完成(仅用于测试,生产环境避免使用block) processTreeDFS(tree).blockLast(); } }
代码解释
thenMany:会等待当前节点的process()处理完成后,才开始订阅并处理子节点的流。flatMap(..., 1):把并发数设为1,确保子节点是按顺序逐个处理的,完美贴合DFS的深度优先逻辑。doOnComplete:只有当thenMany里的所有子节点递归处理流都完成后,才会触发父节点的完成回调,完全符合你的需求。
备选方案:基于expandDeep的调整
如果你一定要用expandDeep,可以通过分组+收集的方式跟踪子节点处理进度,但这种方式需要额外实现父节点查找逻辑,复杂度更高。示例如下:
public static Flux<Void> processTreeWithExpandDeep(Node root) { return Flux.just(root) .expandDeep(node -> Flux.fromIterable(node.getChildren())) // 按父节点分组,需要实现findParentName方法找到每个节点的父节点 .groupBy(node -> findParentName(node, root)) .flatMap(group -> group .flatMap(Node::process) .doOnComplete(() -> { if (group.key() != null) { System.out.println("✅ All children of " + group.key() + " are fully processed!"); } })); } // 辅助方法:递归查找节点的父节点名称 private static String findParentName(Node target, Node currentNode) { for (Node child : currentNode.getChildren()) { if (child.equals(target)) { return currentNode.getName(); } String parentName = findParentName(target, child); if (parentName != null) { return parentName; } } return null; // 根节点A没有父节点 }
这种方式利用groupBy把同一个父节点的所有子节点归为一组,当组内所有节点处理完成时,触发该父节点的完成回调。但缺点是需要额外实现父节点查找逻辑,代码复杂度更高。
总结
最推荐的是递归串联处理的方式,它逻辑清晰,完全贴合DFS的处理顺序,能准确保证父节点的完成回调在所有子节点处理完成后触发,非常适合你当前的场景。
内容的提问来源于stack exchange,提问作者ansh
相关产品推荐
相关产品推荐

