You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

如何基于配置文件编程构建并执行可调用对象组成的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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.10.02 21:24:00