如何基于配置文件编程构建并执行可调用对象组成的DAG
基于配置实现节点可调用对象DAG的构建与执行方案
这个需求本质是实现一个轻量的声明式工作流引擎,核心逻辑是把配置中的依赖关系解析为DAG拓扑结构,再按拓扑序执行对应节点的可调用对象即可,以下是可直接落地的实现思路:
1. 先约定配置文件格式
推荐用YAML做配置,可读性高,配置里需要包含每个执行步骤的唯一标识、目标节点、要调用的可调用对象名称、依赖的上游步骤、上游输出到当前入参的映射规则,你举的例子对应的配置示例如下:
# dag_config.yaml dag: steps: - step_id: step_0 node: "Node1" callable: "A" dependencies: [] input_mapping: {} - step_id: step_1 node: "Node2" callable: "A" dependencies: ["step_0"] input_mapping: {"args": ["{{step_0.output}}"]} - step_id: step_2 node: "Node3" callable: "B" dependencies: ["step_0"] input_mapping: {"args": ["{{step_0.output}}"]} - step_id: step_3 node: "Node1" callable: "B" dependencies: ["step_1", "step_2"] input_mapping: {"args": ["{{step_1.output}}", "{{step_2.output}}"]}
你可以根据业务需求扩展配置字段,比如增加重试次数、超时时间、关键字参数映射等规则。
2. 实现核心逻辑
2.1 节点与可调用对象注册
先实现基础的节点类和注册中心,统一管理所有节点的可调用对象:
from collections import deque import yaml class BaseNode: def __init__(self, node_id): self.node_id = node_id self.callables = {} # 存储当前节点的所有可调用对象,key为名称,value为函数 def get_callable(self, callable_name): if callable_name not in self.callables: raise ValueError(f"节点{self.node_id}未注册名为{callable_name}的可调用对象") return self.callables[callable_name] # 节点注册中心,提前实例化所有固定节点 node_registry = { "Node1": BaseNode("Node1"), "Node2": BaseNode("Node2"), "Node3": BaseNode("Node3") } # 提前给每个节点注册可调用对象,这里是示例,实际可根据业务逻辑初始化 node_registry["Node1"].callables["A"] = lambda: "Node1.A的输出" node_registry["Node1"].callables["B"] = lambda x, y: f"Node1.B的输出,入参为{x}、{y}" node_registry["Node2"].callables["A"] = lambda x: f"Node2.A的输出,入参为{x}" node_registry["Node3"].callables["B"] = lambda x: f"Node3.B的输出,入参为{x}"
2.2 DAG解析与执行器
拓扑排序是整个实现的核心,既可以检测配置中的循环依赖问题,也能保证执行顺序完全符合定义的依赖关系,执行器的完整实现如下:
class DAGExecutor: def __init__(self, config_path): # 加载并解析配置文件 with open(config_path, "r", encoding="utf-8") as f: self.config = yaml.safe_load(f) self.steps = {step["step_id"]: step for step in self.config["dag"]["steps"]} self.step_outputs = {} # 存储所有步骤的执行结果,供下游步骤调用 def _topological_sort(self): # 构建入度表和邻接表 in_degree = {step_id: 0 for step_id in self.steps} adjacency = {step_id: [] for step_id in self.steps} for step_id, step in self.steps.items(): for dep in step["dependencies"]: if dep not in self.steps: raise ValueError(f"步骤{step_id}的依赖{dep}不存在") adjacency[dep].append(step_id) in_degree[step_id] += 1 # 拓扑排序计算执行顺序 queue = deque([sid for sid in in_degree if in_degree[sid] == 0]) topo_order = [] while queue: current_sid = queue.popleft() topo_order.append(current_sid) for next_sid in adjacency[current_sid]: in_degree[next_sid] -= 1 if in_degree[next_sid] == 0: queue.append(next_sid) # 检查是否存在循环依赖 if len(topo_order) != len(self.steps): raise ValueError("配置的DAG存在循环依赖,无法执行") return topo_order def _render_input(self, input_mapping): # 入参模板渲染,把{{step_id.output}}替换为对应步骤的输出 rendered = {"args": [], "kwargs": {}} # 渲染位置参数 for arg in input_mapping.get("args", []): if isinstance(arg, str) and arg.startswith("{{") and arg.endswith("}}"): dep_sid = arg.strip("{}").split(".")[0] rendered["args"].append(self.step_outputs[dep_sid]) else: rendered["args"].append(arg) # 渲染关键字参数 for k, v in input_mapping.get("kwargs", {}).items(): if isinstance(v, str) and v.startswith("{{") and v.endswith("}}"): dep_sid = v.strip("{}").split(".")[0] rendered["kwargs"][k] = self.step_outputs[dep_sid] else: rendered["kwargs"][k] = v return rendered def execute(self): # 按拓扑序执行所有步骤 topo_order = self._topological_sort() for step_id in topo_order: step = self.steps[step_id] # 获取目标节点的可调用对象 node = node_registry[step["node"]] func = node.get_callable(step["callable"]) # 渲染入参 input_params = self._render_input(step["input_mapping"]) # 执行调用并存储结果 output = func(*input_params["args"], **input_params["kwargs"]) self.step_outputs[step_id] = output return self.step_outputs
3. 运行测试
直接初始化执行器调用即可:
if __name__ == "__main__": executor = DAGExecutor("dag_config.yaml") all_outputs = executor.execute() print(all_outputs["step_3"]) # 输出结果为:Node1.B的输出,入参为Node2.A的输出,入参为Node1.A的输出、Node3.B的输出,入参为Node1.A的输出
可选优化点
- 若有并行执行需求,拓扑排序后可以把同一层级无依赖的步骤放入线程池/进程池并行执行,比如示例中的step_1和step_2就可以同时运行
- 需要支持故障重试、执行日志的话,给可调用对象加一层装饰器即可实现
- 入参映射逻辑可以按需扩展,比如支持提取上游输出的指定字段、简单的类型转换等
- 如果不想自行实现全量逻辑,也可以直接用Prefect、Dagster等轻量工作流库,只需把你的节点可调用对象封装成库对应的任务即可。
内容的提问来源于stack exchange,提问作者Rylan Schaeffer
相关产品推荐
相关产品推荐

