news 2026/9/14 12:17:16

DB-GPT AWEL 自定义算子实战:从 MapOperator 到流式算子的完整实现

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
DB-GPT AWEL 自定义算子实战:从 MapOperator 到流式算子的完整实现

DB-GPT AWEL 自定义算子实战:从 MapOperator 到流式算子的完整实现

【免费下载链接】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 的 AWEL(Agentic Workflow Expression Language,智能体工作流表达语言)教程文档《1.3 Custom Operator》为主体,讲解如何在 AWEL DAG 中编写自定义算子:先用MapOperator实现一个最简的“Hello, world!”算子,再用StreamifyAbsOperatorTransformStreamAbsOperator实现一个数字流的生产与转换管道。文中所有示例代码均已在本仓库源码环境下实际运行验证,并给出了对应实现源码(dbgpt/core/awel/operators/)层面的原理剖析,帮助你在 DB-GPT 中开发、调试自己的 LLM 应用工作流节点。

一、AWEL 自定义算子的整体思路

在 AWEL 中,工作流是一个 DAG(有向无环图),图中的每个可执行节点称为算子(Operator)。自定义算子的通用做法只有一条:继承一个基础算子类,重写与之对应的方法。DB-GPT 在 awel 模块入口 中导出了全部基础算子,包括:

基础算子抽象方法作用定义位置
MapOperatorasync def map(self, input: IN) -> OUT一对一数据映射/转换common_operator.py
JoinOperator构造时传入combine_function合并多个上游输入common_operator.py
ReduceStreamOperatorasync def reduce(self, a, b) -> OUT流数据归约为单值common_operator.py
BranchOperatorasync def branches(self)条件分支路由common_operator.py
StreamifyAbsOperatorasync def streamify(self, input: IN) -> AsyncIterator[OUT]单值转异步流stream_operator.py
UnstreamifyAbsOperatorasync def unstreamify(self, it: AsyncIterator[IN]) -> OUT异步流归约为单值stream_operator.py
TransformStreamAbsOperatorasync def transform_stream(self, it: AsyncIterator[IN]) -> AsyncIterator[OUT]流到流的转换stream_operator.py

所有算子的执行入口都在基类 BaseOperator 中定义:

  • call(call_data=..., dag_ctx=..., dag_variables=...)(base.py L262-L286):执行 DAG 并返回最终输出值,非流式场景使用;
  • call_stream(...)(base.py L316-L358):执行 DAG 并返回AsyncIterator,流式场景使用;
  • 两者内部都会调用self._runner.execute_workflow(...),默认 runner 是 DefaultWorkflowRunner。也就是说,无论算子定义得多简单,call都会走完整的 runner → DAG 执行链路,而不仅仅是直接调用你的方法。

二、第一个自定义算子:Hello, world!(MapOperator)

2.1 示例代码

创建一个文件(例如hello_world_custom_operator.py),写入以下代码:

import asyncio from dbgpt.core.awel import DAG, MapOperator class HelloWorldOperator(MapOperator[str, None]): async def map(self, x: str) -> None: print(f"Hello, {x}!") with DAG("awel_hello_world") as dag: task = HelloWorldOperator() asyncio.run(task.call(call_data="world"))

代码要点:

  1. 泛型参数MapOperator[str, None]:第一个参数是输入类型IN,第二个参数是输出类型OUT。这里输出为None,表示算子只产生副作用(打印),不向下游传递数据;
  2. with DAG("awel_hello_world") as dag::进入 DAG 上下文后实例化的算子会自动挂载到该 DAG 上,并自动分配task_id。这一点由元类BaseOperatorMeta_apply_defaults完成——它会在构造算子时自动注入dagtask_idrunnerexecutorvariables_provider等参数(见 base.py L100-L158),因此你在 DAG 上下文里写HelloWorldOperator()不需要传任何参数;
  3. task.call(call_data="world")call_data会作为该“根算子”的输入。注意call内部会把它包装成{"data": "world"}(base.py L280-L281),再通过SimpleCallDataInputSource转成上游输出,最终交给map方法处理。

2.2 运行与验证

在仓库环境下运行:

python hello_world_custom_operator.py

输出(本仓库环境实测一致):

Hello, world!

2.3 源码视角:map是如何被调用的

MapOperator的执行逻辑在_do_run中(common_operator.py L155-L198),核心流程:

  1. 取出当前TaskContextcall_data
  2. 如果存在call_data(本例即world),先校验父节点数量约束:MapOperator要求“单父节点”或者由call_data直接驱动(L170-L177,否则抛出MapDAGNode expects single parent错误);
  3. 选择映射函数:map_function = self.map_function or self.map。也就是说,MapOperator支持两种等价写法——继承类并覆盖map方法(本例),或构造时直接传入一个可调用对象MapOperator(map_function=...)
  4. 调用wrapped_call_data.map(map_function)得到新的TaskOutput,并set_task_output写回上下文,供下游算子消费。

另外,若你通过构造函数传入map_function,且算子启用了序列化检查(check_serializable),AWEL 会验证该函数是否可 pickle(L148-L151),这是为了支持工作流持久化与远程执行场景。

三、第一个流式算子:数字流生产 + 流转换

3.1 示例代码

流式(streaming)是 AWEL 的一等公民能力:一个算子的输出可以是AsyncIterator,下游算子可以在数据“产生即消费”的方式上逐条处理,天然适合 LLM 增量输出、大列表分页处理等场景。示例创建custom_streaming_operator.py

import asyncio from typing import AsyncIterator from dbgpt.core.awel import DAG, StreamifyAbsOperator, TransformStreamAbsOperator class NumberProducerOperator(StreamifyAbsOperator[int, int]): async def streamify(self, n: int) -> AsyncIterator[int]: for i in range(n): yield i class NumberDoubleOperator(TransformStreamAbsOperator[int, int]): async def transform_stream(self, it: AsyncIterator) -> AsyncIterator[int]: async for i in it: # Double the number yield i * 2 with DAG("numbers_dag") as dag: task = NumberProducerOperator() double_task = NumberDoubleOperator() task >> double_task async def helper_call_fn(t, n: int): # Call the streaming operator by `call_stream` method async for i in await t.call_stream(call_data=n): print(i) asyncio.run(helper_call_fn(double_task, 10))

代码要点:

  1. NumberProducerOperatorstreamify(self, n)把单个整数n转成0..n-1的异步迭代器,yield使得数据按需逐个产生;
  2. NumberDoubleOperatortransform_stream接收上游的AsyncIterator,逐条取出、翻倍后再yield,实现“流到流”的转换,全程不需要把整条流加载进内存;
  3. task >> double_task:AWEL 复用 Python 的>>运算符定义 DAG 边的语法糖,把生产算子的输出接到转换算子的输入;
  4. await t.call_stream(call_data=n):调用链是call_streamrunner.execute_workflow(streaming_call=True)→ 返回AsyncIterator,再配合async for逐条消费。

3.2 运行与验证

python custom_streaming_operator.py

输出(本仓库环境实测一致):

0 2 4 6 8 10 12 14 16 18

3.3 源码视角:流式算子与普通算子的区别

  • StreamifyAbsOperatorTransformStreamAbsOperator都声明了类属性streaming_operator = True(stream_operator.py L14、L89),AWEL 据此识别流式节点;
  • StreamifyAbsOperator._do_run(L16-L32)把call_data包装后调用task_output.streamify(self.streamify),产出的是SimpleStreamTaskOutput
  • 消费端的call_stream(base.py L316-L358)有两个值得注意的健壮性设计:
    • 它检查task_output.is_stream;如果最终输出不是流,会自动把单个值包装成一个只yield一次的生成器(L350-L355)。这意味着对非流式算子调用call_stream也不会报错,只是得到单元素流——但教程中强调:对真正的流式算子,必须await拿到迭代器后再async for遍历,这是初学者最容易踩的坑;
    • 整个迭代过程被root_tracer.wrapper_async_stream包了一层 span(dbgpt.awel.operator.call_stream.iterate),便于链路追踪排查流式断流问题。

四、实践建议:如何选择基础算子、如何调试

结合上面的源码结构,给出几条可直接落地的建议:

  1. 按数据形态选基类:输入输出都是单值 →MapOperator;输入单值、输出多条 →StreamifyAbsOperator;输入输出都是流 →TransformStreamAbsOperator;输入流、输出单值 →UnstreamifyAbsOperator;多个上游汇聚 →JoinOperator;条件路由 →BranchOperator。各类型的输入约束(如ReduceStreamOperator要求“流 + 单父节点”,见 common_operator.py L110-L113)在_do_run中都有显式校验,违反时会抛出带节点信息的ValueError,可直接据此定位 DAG 连接错误;
  2. call_data的传递规则call_data只作用于 DAG 的根算子。在示例中task.call(call_data="world")"world"进入的是HelloWorldOperator.mapcall_stream(call_data=10)10进入的是NumberProducerOperator.streamify。中间节点之间靠 DAG 边传递数据,而不是call_data
  3. 阻塞代码要放进执行器:AWEL 的算子全部在事件循环中异步执行,如果你的算子里有同步阻塞调用(如传统数据库驱动),应使用基类提供的blocking_func_to_async(func, ...)(base.py L389-L404),它会把阻塞函数丢到self._executor线程池执行,避免卡死整个事件循环;
  4. 开发期调试:在纯开发场景下可以直接asyncio.run(task.call(...))(本文两个示例即如此);如果 DAG 中带HttpTrigger等触发器,仓库提供了 setup_dev_environment,可一键启动本地 HTTP 服务(默认127.0.0.1:5555)、注册触发器并可视化 DAG 图(依赖 graphviz),适合联调含触发器的完整工作流;
  5. 变量占位符:基类在执行前会调用_resolve_variables(base.py L415-L490),把算子属性中的VariablesPlaceHolder按“DAG 变量优先、系统变量兜底”的顺序解析。从源码结构看,自定义算子若把提示词、模型名等配置写成变量占位符,即可在运行时被动态注入,而无需修改算子代码。

五、小结

本文完整复现并验证了 AWEL 教程《1.3 Custom Operator》的两个核心示例:继承MapOperator覆盖map方法得到最简自定义算子,继承StreamifyAbsOperator/TransformStreamAbsOperator并用>>连接得到流式管道,再用call/call_stream两种入口驱动执行。其背后的统一执行链路是:call/call_streamDefaultWorkflowRunner.execute_workflow→ 各算子_do_runmap/streamify/transform_stream)→TaskOutput沿 DAG 边传递。掌握这套“继承 + 重写 + DAG 编排”的模式后,你就可以在 DB-GPT 的 Agent、RAG、HTTP 服务等上层功能中,自由编写任意业务节点。

【免费下载链接】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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/14 12:16:35

html2canvas+jsPDF PDF截断的像素级定位与分页修复

简介:本资源聚焦前端 PDF 生成场景中 html2canvas 与 jsPDF 结合使用时常见的内容截断难题,面向 Web 开发者、前端工程师及需要导出长页面为 PDF 的项目实践者。方案通过创新的像素级扫描逻辑识别截断位置:先将 HTML 渲染为白色背景图片&…

作者头像 李华
网站建设 2026/9/14 12:16:18

GitHub Copilot替代方案全解析:从免费工具到付费IDE横向评测

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/14 12:15:49

企业AI效能管理:从模型上线到持续治理的落地指南

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/14 12:15:46

.NET Framework 4.6.1 电商源码部署指南:Himall3.0 商城实战配置

简介:本资源为Himall3.0电子商务平台完整开源源码包,面向Java/Python/Node.js等技术栈的中高级开发者、电商系统学习者及二次开发需求者,提供可研究、可定制、可部署的成熟商城系统实践样本。压缩包大小376.2MB,虽未提供具体文件总…

作者头像 李华
网站建设 2026/9/14 12:14:36

AI Agent技术解析与实战:从架构到应用

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华