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!”算子,再用StreamifyAbsOperator与TransformStreamAbsOperator实现一个数字流的生产与转换管道。文中所有示例代码均已在本仓库源码环境下实际运行验证,并给出了对应实现源码(dbgpt/core/awel/operators/)层面的原理剖析,帮助你在 DB-GPT 中开发、调试自己的 LLM 应用工作流节点。
一、AWEL 自定义算子的整体思路
在 AWEL 中,工作流是一个 DAG(有向无环图),图中的每个可执行节点称为算子(Operator)。自定义算子的通用做法只有一条:继承一个基础算子类,重写与之对应的方法。DB-GPT 在 awel 模块入口 中导出了全部基础算子,包括:
| 基础算子 | 抽象方法 | 作用 | 定义位置 |
|---|---|---|---|
MapOperator | async def map(self, input: IN) -> OUT | 一对一数据映射/转换 | common_operator.py |
JoinOperator | 构造时传入combine_function | 合并多个上游输入 | common_operator.py |
ReduceStreamOperator | async def reduce(self, a, b) -> OUT | 流数据归约为单值 | common_operator.py |
BranchOperator | async def branches(self) | 条件分支路由 | common_operator.py |
StreamifyAbsOperator | async def streamify(self, input: IN) -> AsyncIterator[OUT] | 单值转异步流 | stream_operator.py |
UnstreamifyAbsOperator | async def unstreamify(self, it: AsyncIterator[IN]) -> OUT | 异步流归约为单值 | stream_operator.py |
TransformStreamAbsOperator | async 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"))代码要点:
- 泛型参数
MapOperator[str, None]:第一个参数是输入类型IN,第二个参数是输出类型OUT。这里输出为None,表示算子只产生副作用(打印),不向下游传递数据; with DAG("awel_hello_world") as dag::进入 DAG 上下文后实例化的算子会自动挂载到该 DAG 上,并自动分配task_id。这一点由元类BaseOperatorMeta的_apply_defaults完成——它会在构造算子时自动注入dag、task_id、runner、executor、variables_provider等参数(见 base.py L100-L158),因此你在 DAG 上下文里写HelloWorldOperator()不需要传任何参数;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),核心流程:
- 取出当前
TaskContext的call_data; - 如果存在
call_data(本例即world),先校验父节点数量约束:MapOperator要求“单父节点”或者由call_data直接驱动(L170-L177,否则抛出MapDAGNode expects single parent错误); - 选择映射函数:
map_function = self.map_function or self.map。也就是说,MapOperator支持两种等价写法——继承类并覆盖map方法(本例),或构造时直接传入一个可调用对象MapOperator(map_function=...); - 调用
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))代码要点:
NumberProducerOperator:streamify(self, n)把单个整数n转成0..n-1的异步迭代器,yield使得数据按需逐个产生;NumberDoubleOperator:transform_stream接收上游的AsyncIterator,逐条取出、翻倍后再yield,实现“流到流”的转换,全程不需要把整条流加载进内存;task >> double_task:AWEL 复用 Python 的>>运算符定义 DAG 边的语法糖,把生产算子的输出接到转换算子的输入;await t.call_stream(call_data=n):调用链是call_stream→runner.execute_workflow(streaming_call=True)→ 返回AsyncIterator,再配合async for逐条消费。
3.2 运行与验证
python custom_streaming_operator.py输出(本仓库环境实测一致):
0 2 4 6 8 10 12 14 16 183.3 源码视角:流式算子与普通算子的区别
StreamifyAbsOperator和TransformStreamAbsOperator都声明了类属性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),便于链路追踪排查流式断流问题。
- 它检查
四、实践建议:如何选择基础算子、如何调试
结合上面的源码结构,给出几条可直接落地的建议:
- 按数据形态选基类:输入输出都是单值 →
MapOperator;输入单值、输出多条 →StreamifyAbsOperator;输入输出都是流 →TransformStreamAbsOperator;输入流、输出单值 →UnstreamifyAbsOperator;多个上游汇聚 →JoinOperator;条件路由 →BranchOperator。各类型的输入约束(如ReduceStreamOperator要求“流 + 单父节点”,见 common_operator.py L110-L113)在_do_run中都有显式校验,违反时会抛出带节点信息的ValueError,可直接据此定位 DAG 连接错误; call_data的传递规则:call_data只作用于 DAG 的根算子。在示例中task.call(call_data="world")的"world"进入的是HelloWorldOperator.map;call_stream(call_data=10)的10进入的是NumberProducerOperator.streamify。中间节点之间靠 DAG 边传递数据,而不是call_data;- 阻塞代码要放进执行器:AWEL 的算子全部在事件循环中异步执行,如果你的算子里有同步阻塞调用(如传统数据库驱动),应使用基类提供的
blocking_func_to_async(func, ...)(base.py L389-L404),它会把阻塞函数丢到self._executor线程池执行,避免卡死整个事件循环; - 开发期调试:在纯开发场景下可以直接
asyncio.run(task.call(...))(本文两个示例即如此);如果 DAG 中带HttpTrigger等触发器,仓库提供了 setup_dev_environment,可一键启动本地 HTTP 服务(默认127.0.0.1:5555)、注册触发器并可视化 DAG 图(依赖 graphviz),适合联调含触发器的完整工作流; - 变量占位符:基类在执行前会调用
_resolve_variables(base.py L415-L490),把算子属性中的VariablesPlaceHolder按“DAG 变量优先、系统变量兜底”的顺序解析。从源码结构看,自定义算子若把提示词、模型名等配置写成变量占位符,即可在运行时被动态注入,而无需修改算子代码。
五、小结
本文完整复现并验证了 AWEL 教程《1.3 Custom Operator》的两个核心示例:继承MapOperator覆盖map方法得到最简自定义算子,继承StreamifyAbsOperator/TransformStreamAbsOperator并用>>连接得到流式管道,再用call/call_stream两种入口驱动执行。其背后的统一执行链路是:call/call_stream→DefaultWorkflowRunner.execute_workflow→ 各算子_do_run(map/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),仅供参考