带执行依赖约束的流水线图遍历算法伪代码求解
DataFrame运算流水线依赖有序遍历方案
已知约束梳理
- 流水线节点共3类输入输出形态:
- 单输入双输出(one-input-->two-outputs)
- 单输入单输出(one-input-->one-output)
- 双输入单输出(two-inputs-->one-output)
- 图为从左到右单向流转的类DAG结构,无环,复杂度低于通用DAG
- 已知全部起始节点为
[S1, S2, S3, S4],存在明确终端节点 - 节点执行有严格前置依赖要求:例如执行
operation2前必须完成operation1与operation5,执行operation7前必须完成operation3与operation6
核心实现思路
因为结构比通用DAG简单,不需要提前做全量拓扑排序计算,用入度计数+动态就绪队列的轻量逻辑即可覆盖需求:
- 为每个节点维护三个基础属性:前置依赖总数、已满足的依赖计数、下游节点列表
- 初始将所有无前置依赖的起始节点加入就绪执行队列
- 每次从队列取出可执行节点运行,运行完成后遍历它的所有下游节点,给下游的已满足依赖计数+1
- 只要下游节点的已满足依赖计数等于它需要的总前置依赖数,说明该节点已经具备执行条件,加入就绪队列
- 重复调度直到终端节点执行完成,最后做全节点执行校验即可
伪代码实现
// 第一步:初始化所有节点的基础属性 FOR each node in 全量节点集合: node.in_degree = 节点前置依赖列表长度 // 执行该节点需要的前置依赖总数 node.satisfied_dep_count = 0 // 已经执行完成的前置依赖数量 node.downstream_nodes = [] // 存储该节点直接指向的所有下游节点 END FOR // 反向填充所有节点的下游关系 FOR each node in 全量节点集合: FOR each dep in node.pre_dependencies: dep.downstream_nodes.append(node) END FOR END FOR // 初始化就绪队列,塞入所有起始节点 ready_queue = 先进先出队列() FOR start_node in [S1, S2, S3, S4]: ready_queue.push(start_node) END FOR executed_set = 空集合 // 记录已经执行完成的节点 // 开始按依赖顺序调度执行 WHILE ready_queue 不为空: current_node = ready_queue.pop() // 执行当前节点绑定的DataFrame操作逻辑 执行节点运算(current_node) executed_set.add(current_node) // 运行到终端节点即可终止调度 IF current_node == 终端节点: BREAK END IF // 更新所有下游节点的依赖满足状态 FOR downstream_node in current_node.downstream_nodes: downstream_node.satisfied_dep_count += 1 // 下游所有前置依赖都凑齐了,加入就绪队列等待执行 IF downstream_node.satisfied_dep_count == downstream_node.in_degree: ready_queue.push(downstream_node) END IF END FOR END WHILE // 最终校验,避免存在依赖缺失导致的节点漏执行 IF executed_set长度 != 全量节点集合长度: 抛出异常("存在依赖无法满足的节点,流水线执行失败") END IF
三类节点的适配说明
- 单输入单输出节点:入度固定为1,唯一上游执行完成后会直接触发该节点入队,完全匹配一对一的流转逻辑
- 单输入双输出节点:执行完成后会同步更新两个下游节点的依赖计数,两个下游各自等待自身剩余依赖满足后即可执行,适配分流场景
- 双输入单输出节点:入度固定为2,必须等两个上游全部执行完成、依赖计数累计到2才会入队,天然适配合流场景需要等两个输入都就绪的约束
扩展提示:如果后续需要加分支跳过逻辑,只需要在当前节点执行完成后判断哪些下游需要触发,不给跳过的下游更新依赖计数即可,不需要改动核心调度逻辑。
内容的提问来源于stack exchange,提问作者Ritik Kamra
相关产品推荐
相关产品推荐

