实战指南)
DB-GPT AWEL 任务生命周期钩子Task Lifecycle Hooks实战指南【免费下载链接】DB-GPTopen-source agentic AI data assistant for the next generation of AI Data products.项目地址: https://gitcode.com/GitHub_Trending/db/DB-GPT本文以 DB-GPT 的 AWELAgentic Workflow Expression LanguageAgent 工作流表达式语言框架为核心系统讲解任务生命周期钩子before_dag_run与after_dag_end的概念、完整代码写法、运行方式及其在源码层面的执行时机与调用链。读完本文你将掌握在任意 AWEL 任务节点中嵌入「DAG 启动前 / 结束后」自定义逻辑的能力并理解流式与批式场景下生命周期回调的差异可直接用于资源初始化、指标统计、上下文清理等真实工作流场景。一、什么是 AWEL 任务生命周期钩子AWEL 是 DB-GPT 内置的工作流编排框架其核心抽象是DAG有向无环图与挂载在 DAG 上的各类算子Operator。一个典型 AWEL 流程由触发节点Trigger、处理节点Operator如MapOperator和结束节点构成数据沿着 DAG 边流动。在 DAG 的一次完整执行过程中除了节点自身的业务逻辑例如MapOperator的map()方法之外框架还预留了一组任务生命周期钩子Task Lifecycle Hooks它们是可在任务中实现的一组方法用于在任务生命周期的不同阶段执行动作。AWEL 目前提供两个钩子before_dag_run在 DAG 开始运行之前执行after_dag_end在 DAG 结束之后执行。这两个钩子定义在源码层的DAGLifecycle基类中packages/dbgpt-core/src/dbgpt/core/awel/dag/base.py#L282-L294class DAGLifecycle: The lifecycle of DAG. async def before_dag_run(self): Execute before DAG run. pass async def after_dag_end(self, event_loop_task_id: int): Execute after DAG end. This method may be called multiple times, please make sure it is idempotent. pass值得注意的是after_dag_end的源码注释明确提示该方法可能被多次调用请确保其实现是幂等的。这一点在编写清理型逻辑时尤其重要详见本文第六节。在继承体系中DAGNode是DAGLifecycle的子类packages/dbgpt-core/src/dbgpt/core/awel/dag/base.py#L297而所有算子节点Operator又都继承自DAGNode因此任何自定义算子都可以直接覆写这两个钩子。二、核心概念两个钩子的职责边界before_dag_run—— DAG 运行前触发时机DefaultWorkflowRunner.execute_workflow()在正式执行 DAG 节点之前调用见 packages/dbgpt-core/src/dbgpt/core/awel/runner/local_runner.py#L102。典型用途连接池预热、模型/向量索引加载、配置校验、全局参数预置、日志标记等一次性初始化工作。特点对整张 DAG 的所有节点统一触发且在数据真正开始在节点间流转之前执行因此适合放置「先决条件」类逻辑。after_dag_end—— DAG 结束后触发时机DAG 主体执行完毕、节点输出都已落盘后触发。以默认本地执行器为例在execute_workflow()的非流式分支中执行完_execute_node后会调用node.dag._after_dag_end(...)见 packages/dbgpt-core/src/dbgpt/core/awel/runner/local_runner.py#L117-L120。典型用途释放连接与资源、上报耗时指标、持久化最终结果、发送通知、清理临时状态等收尾工作。特点接收一个event_loop_task_id参数当前事件循环任务的唯一标识默认取id(asyncio.current_task())用于在并发的多次 DAG 执行中区分不同的运行实例。三、完整实战示例第一个生命周期钩子任务原教程要求在awel_tutorial目录下新建lifecycle_hooks.py文件并写入以下代码。该文件位于英文版教程目录docs/docs/awel/awel_tutorial/lifecycle_hooks.py你也可以在本地 DB-GPT 源码目录中自行创建import asyncio from dbgpt.core.awel import DAG, MapOperator class MyLifecycleTask(MapOperator[str, str]): async def before_dag_run(self): print(Before DAG run) async def after_dag_end(self): print(After DAG end) async def map(self, x: str) - str: return fHello, {x}! with DAG(awel_lifecycle_hooks) as dag: task MyLifecycleTask() print(asyncio.run(task.call(world)))下面逐段拆解这段代码的关键点导入dbgpt.core.awel聚合导出了DAG与MapOperator它们是编写 AWEL 流程的最小依赖。自定义任务MyLifecycleTask继承自MapOperator[str, str]即输入输出均为str的映射算子。它覆写了两个生命周期钩子并实现了核心业务方法map()。从源码看MapOperator的_do_run()会优先使用构造时传入的map_function否则回退到子类覆写的map()方法packages/dbgpt-core/src/dbgpt/core/awel/operators/common_operator.py#L178-L198本例即采用「覆写map()」的方式。DAG 上下文with DAG(awel_lifecycle_hooks) as dag:创建了一个名为awel_lifecycle_hooks的 DAG并将task挂载其中离开with块时 DAG 自动完成构建。调用task.call(world)从结束节点视角驱动整条工作流执行asyncio.run(...)负责驱动事件循环。四、运行方式与预期输出在原文档的示例场景中运行命令为poetry run python awel_tutorial/lifecycle_hooks.py即通过 Poetry 管理的虚拟环境运行脚本。如果你不依赖 Poetry也可以直接使用本机 Python 解释器执行python awel_tutorial/lifecycle_hooks.py前提是当前环境中已安装 DB-GPT 及其 AWEL 依赖参见仓库根目录的pyproject.toml与uv.lock。若你的脚本放在其他位置将上述相对路径替换为实际路径即可。预期输出如下Before DAG run After DAG end Hello, world!输出顺序正好印证了两个钩子的触发语义Before DAG run出现在业务结果之前After DAG end出现在业务结果之后call()返回前完成收尾。五、源码级原理钩子的执行时机与调用链要真正理解生命周期钩子需要沿着「调用入口 → 执行器 → 调度器 → 节点」这条链路看一遍源码。5.1 入口默认工作流执行器所有 AWEL 工作流的默认执行器是DefaultWorkflowRunnerpackages/dbgpt-core/src/dbgpt/core/awel/runner/local_runner.py#L24。在execute_workflow()中钩子的调用顺序非常清晰通过JobManager.build_from_end_node(node, call_data)从结束节点反向构建整张 DAG 的节点集合与调用数据创建并保存DAGContext记录event_loop_task_id、节点输出、共享数据等调用await job_manager.before_dag_run()——全节点 before 钩子在root_tracer的追踪跨度内执行_execute_node(...)递归完成所有上游节点到结束节点的实际计算非流式且非子 DAG 场景下调用await node.dag._after_dag_end(dag_ctx._event_loop_task_id)——全节点 after 钩子packages/dbgpt-core/src/dbgpt/core/awel/runner/local_runner.py#L102-L120。5.2 调度JobManager 的广播式调用JobManager本身也继承自DAGLifecyclepackages/dbgpt-core/src/dbgpt/core/awel/runner/job_manager.py#L14其实现是对整张 DAG 所有节点的广播调用并使用asyncio.gather并行执行避免节点数较多时串行等待async def before_dag_run(self): Execute the callback before DAG run. tasks [] for node in self._all_nodes: tasks.append(node.before_dag_run()) await asyncio.gather(*tasks) async def after_dag_end(self, event_loop_task_id: int): Execute the callback after DAG end. tasks [] for node in self._all_nodes: tasks.append(node.after_dag_end(event_loop_task_id)) await asyncio.gather(*tasks)见 packages/dbgpt-core/src/dbgpt/core/awel/runner/job_manager.py#L75-L87也就是说只要 DAG 中的任意一个节点覆写了钩子都会在对应时机被并行触发没有覆写的节点使用基类空实现直接跳过。5.3 流式场景after 钩子如何被延迟到流结束对于流式调用streaming_callTrueexecute_workflow()不会在返回DAGContext前同步调用_after_dag_end。以 HTTP 触发场景为例HttpTrigger会把_after_dag_end注册为BackgroundTasks回调待StreamingResponse的生成器消费完毕后由 FastAPI 的响应后台任务统一执行packages/dbgpt-core/src/dbgpt/core/awel/trigger/http_trigger.py#L746-L750async def _after_dag_end(): await dag._after_dag_end(end_node.current_event_loop_task_id) background_tasks BackgroundTasks() background_tasks.add_task(_after_dag_end)同理迭代器触发器IteratorTrigger在流式消费结束后也会显式调用dag._after_dag_end(...)packages/dbgpt-core/src/dbgpt/core/awel/trigger/iterator_trigger.py#L233。因此可以推断在流式工作流中after_dag_end的触发点是「整个响应流完整结束」而不是「首个输出返回」编写清理逻辑时需要把这一点考虑在内。5.4 DAG 层收口与上下文清理DAG._after_dag_end()是所有触发场景的最终收口实现它遍历self.node_map中的全部节点并行调用after_dag_end随后按event_loop_task_id清理_event_loop_task_id_to_ctx中的 DAG 上下文packages/dbgpt-core/src/dbgpt/core/awel/dag/base.py#L918-L939。这意味着after_dag_end不只负责用户自定义收尾还参与了框架自身的运行时上下文回收是每次 DAG 执行必经的生命周期环节。六、典型应用场景与编写建议基于钩子的语义和源码约束以下是几种贴近实战的用法场景建议使用的钩子说明加载模型 / 建立连接池 / 初始化全局配置before_dag_run在数据流动前一次性完成避免在map()内重复初始化释放连接 / 关闭客户端 / 清理临时目录after_dag_end无论业务是否抛出异常都应保证幂等与可重入记录 DAG 总耗时、上报指标after_dag_end此时所有节点已完成可统计完整执行周期为整次运行写入追踪标记before_dag_run可通过DAGContext/ 共享数据向下游传递编写时请务必遵守两条源码级约束幂等性after_dag_end可能被多次调用源码 docstring 明确要求幂等。释放资源前应判断是否已释放或使用标志位保护否则在流式、重试或子 DAG 场景下可能出现重复执行。轻量快速before_dag_run与after_dag_end在JobManager中是全节点并行调用的如果单个钩子内存在耗时 I/O 或阻塞调用会拉长整个 DAG 的总执行时间建议只放必要的初始化与收尾逻辑。七、更进一步在真实 AWEL 流程中组合使用生命周期钩子并非只能在单独的任务类中使用它可以与任何 AWEL 算子组合。参考 示例脚本 与 数据助理示例一个常见的组合模式是用before_dag_run准备会话级上下文用MapOperator处理用户输入再用after_dag_end汇总输出并清理资源。由于钩子定义在算子基类上基于DAGNode派生的所有节点包括自定义算子、工具节点等都能无侵入地获得这一能力这也是将横切逻辑初始化、清理、观测从业务算子中解耦出来的推荐实践。八、小结AWEL 提供了before_dag_run与after_dag_end两个任务生命周期钩子分别对应 DAG 运行前与 DAG 结束后的执行时机只需继承任意算子如MapOperator并覆写钩子方法即可零侵入地挂载初始化与收尾逻辑源码层面钩子由DefaultWorkflowRunner驱动、JobManager全节点并行广播、DAG._after_dag_end统一收口流式场景下after_dag_end延迟到响应流完全结束后触发编写钩子时务必保证幂等与轻量并善用event_loop_task_id区分并发运行实例。如果你希望进一步理解 AWEL 的算子体系MapOperator、BranchOperator等与 DAG 构建细节可继续阅读仓库中的 AWEL 教程 与 AWEL 文档。【免费下载链接】DB-GPTopen-source agentic AI data assistant for the next generation of AI Data products.项目地址: https://gitcode.com/GitHub_Trending/db/DB-GPT创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考