DB-GPT AWEL 教程:用 UnstreamifyAbsOperator 将数据流聚合为单个结果
【免费下载链接】DB-GPTopen-source agentic AI data assistant for the next generation of AI + Data products.项目地址: https://gitcode.com/GitHub_Trending/db/DB-GPT
UnstreamifyAbsOperator是 DB-GPT 的 AWEL(Agentic Workflow Expression Language,Agent 工作流表达式语言)中负责"流收束"的抽象算子,与StreamifyAbsOperator恰好相反:前者把一条数据展开成异步数据流,后者则把异步数据流收敛成单个数据。本篇文章以官方教程 2.6_unstreamify_operator.md 为骨架,结合仓库源码逐层拆解其用法与原理。读完本文,你将掌握自定义UnstreamifyAbsOperator的完整套路,能够在 DAG 中把流式输出(如流式 LLM 的 token 序列)聚合为最终结果,并理解 AWEL 流算子的底层执行链路。
UnstreamifyAbsOperator 是什么
在 AWEL 中,算子是构成 DAG(有向无环图)的最小执行单元。流式场景下共有三类与"流"直接相关的抽象算子,它们共同定义在 packages/dbgpt-core/src/dbgpt/core/awel/operators/stream_operator.py 中:
| 抽象算子 | 方向 | 核心抽象方法 | 典型用途 |
|---|---|---|---|
StreamifyAbsOperator | 单数据 → 数据流 | streamify(input_value: IN) -> AsyncIterator[OUT] | 把单次输入展开成多条输出,如模拟流式 LLM 服务逐字返回 |
UnstreamifyAbsOperator | 数据流 → 单数据 | unstreamify(input_value: AsyncIterator[IN]) -> OUT | 把上游流式输出聚合为单个最终结果,如对数字流求和 |
TransformStreamAbsOperator | 数据流 → 数据流 | transform_stream(input_value: AsyncIterator[IN]) -> AsyncIterator[OUT] | 流到流的逐条变换,如对每个元素加一 |
从类签名上看,UnstreamifyAbsOperator继承自BaseOperator[OUT]并带有两个泛型参数IN与OUT(stream_operator.py):
class UnstreamifyAbsOperator(BaseOperator[OUT], Generic[IN, OUT]):其中IN是上游传入数据流中每个元素的类型,OUT是聚合后返回的单个结果的类型。它唯一的抽象方法声明为:
@abstractmethod async def unstreamify(self, input_value: AsyncIterator[IN]) -> OUT: """Convert a value of AsyncIterator[IN] to an OUT."""与 Reduce 算子的区别
AWEL 还提供了通用的MapOperator、ReduceOperator等算子(参见 2.2_reduce_operator.md)。UnstreamifyAbsOperator与它们的本质区别在于输入形态:普通算子拿到的是父节点输出的一条具体数据,而UnstreamifyAbsOperator的输入是AsyncIterator[IN],即必须通过async for逐条消费的异步迭代器。因此它天然适合接在StreamifyAbsOperator或TransformStreamAbsOperator等流式节点之后,作为整个流式链路的"终点聚合器"。
自定义一个 UnstreamifyAbsOperator
按照官方教程,使用方式只有一种:继承UnstreamifyAbsOperator并覆写unstreamify方法,在该方法中消费异步迭代器并返回单个结果。下面是最小可用的自定义算子:
from typing import AsyncIterator from dbgpt.core.awel import DAG, UnstreamifyAbsOperator class SumOperator(UnstreamifyAbsOperator[int, int]): """Unstreamify the stream of numbers""" async def unstreamify(self, it: AsyncIterator[int]) -> int: return sum([i async for i in it]) with DAG("sum_dag") as dag: task = SumOperator()这段代码做了三件事:
- 通过
from dbgpt.core.awel import ...引入DAG与UnstreamifyAbsOperator。这两个符号在 packages/dbgpt-core/src/dbgpt/core/awel/__init__.py 中均被显式导出(UnstreamifyAbsOperator同时出现在__all__列表中)。 - 定义
SumOperator,泛型参数[int, int]表示"输入数据流元素为 int、聚合结果为 int"。 - 覆写
unstreamify,用异步列表推导式sum([i async for i in it])一次性消费整个流并求和。
覆写 unstreamify 时需要注意什么
源码 docstring 给出了另一个常见实现(stream_operator.py):
class MyUnstreamOperator(UnstreamifyAbsOperator[int, int]): async def unstreamify(self, input_value: AsyncIterator[int]) -> int: value_cnt = 0 async for v in input_value: value_cnt += 1 return value_cnt由此可以总结出三点实战要点:
- 必须用
async for消费:input_value是异步迭代器,不能直接用普通for遍历; - 返回值必须是单个值而非迭代器:
unstreamify的职责就是把流"关掉",返回OUT类型的普通数据; - 泛型要贴合实际:
IN描述流的元素类型,OUT描述返回值类型,类型标注错误不会报错但会误导阅读者与 IDE 推断。
底层原理:UnstreamifyAbsOperator 如何执行
UnstreamifyAbsOperator并不需要你自己实现执行逻辑,AWEL 框架在其_do_run中完成了"取上游输出 → 应用 unstreamify → 写回结果"的全过程(stream_operator.py):
async def _do_run(self, dag_ctx: DAGContext) -> TaskOutput[OUT]: curr_task_ctx: TaskContext[OUT] = dag_ctx.current_task_context output: TaskOutput[OUT] = await curr_task_ctx.task_input.parent_outputs[ 0 ].task_output.unstreamify(self.unstreamify) curr_task_ctx.set_task_output(output) return output这条链路包含三个关键环节:
- 取上游输出:通过
curr_task_ctx.task_input.parent_outputs[0].task_output拿到父节点的TaskOutput。注意它固定取第一个父输出,因此UnstreamifyAbsOperator是单输入算子,DAG 中只能有一个上游节点直接连接到它; - 调用 TaskOutput.unstreamify:真正的工作发生在 packages/dbgpt-core/src/dbgpt/core/awel/task/task_impl.py 的
SimpleStreamTaskOutput.unstreamify中:
async def unstreamify(self, transform_func: UnStreamFunc) -> TaskOutput[OUT]: if asyncio.iscoroutinefunction(transform_func): out = await transform_func(self.output_stream) else: out = transform_func(self.output_stream) return SimpleTaskOutput(out)可以看到它把self.output_stream(即AsyncIterator[IN])传给你覆写的unstreamify方法,再把返回值包装成一个非流式的SimpleTaskOutput。UnStreamFunc的类型别名定义在 packages/dbgpt-core/src/dbgpt/core/awel/task/base.py:
UnStreamFunc = Callable[[AsyncIterator[IN]], OUT]- 写回任务上下文:聚合出的
SimpleTaskOutput通过set_task_output写入当前任务上下文,供下游节点或最终调用方读取。
值得注意的是:UnstreamifyAbsOperator覆写的unstreamify必须声明为async(源码类型定义为协程函数签名async def unstreamify(...)),即使实现体内不await任何东西也要保留async关键字,因为_do_run期望拿到的是一个可等待对象。
完整实战:对数字流求和
官方教程给出了一个端到端可运行的完整示例:用StreamifyAbsOperator产出0..n-1的数字流,再用SumOperator聚合求和。在awel_tutorial目录下新建文件unstreamify_operator_sum_numbers.py,内容如下:
import asyncio from typing import AsyncIterator from dbgpt.core.awel import DAG, UnstreamifyAbsOperator, StreamifyAbsOperator class NumberProducerOperator(StreamifyAbsOperator[int, int]): """Create a stream of numbers from 0 to `n-1`""" async def streamify(self, n: int) -> AsyncIterator[int]: for i in range(n): yield i class SumOperator(UnstreamifyAbsOperator[int, int]): """Unstreamify the stream of numbers""" async def unstreamify(self, it: AsyncIterator[int]) -> int: return sum([i async for i in it]) with DAG("sum_dag") as dag: task = NumberProducerOperator() sum_task = SumOperator() task >> sum_task print(asyncio.run(sum_task.call(call_data=5))) print(asyncio.run(sum_task.call(call_data=10)))这段示例里有几个值得展开的细节:
task >> sum_task是 DAG 连线语法:表示NumberProducerOperator的输出流入SumOperator,二者构成一条两节点的流水线。教程配套的 2.1_map_operator.md 等章节对连线语法有更系统的介绍;- 调用终点是
SumOperator:sum_task.call(call_data=5)从 DAG 的根节点(NumberProducerOperator)注入call_data=5,工作流沿连线流动,最终在SumOperator处完成聚合并返回int; call返回的是最终值而非流:call方法定义在 packages/dbgpt-core/src/dbgpt/core/awel/operators/base.py,它会执行完整工作流并返回out_ctx.current_task_context.task_output.output。与流式场景下的call_stream(base.py)不同,call直接把内部聚合结果暴露出来,这正是UnstreamifyAbsOperator适合作为"流式链路终点"的原因——下游调用方无需再消费流,直接拿到最终答案。
运行与预期输出
在仓库根目录执行:
poetry run python awel_tutorial/unstreamify_operator_sum_numbers.py控制台将打印:
10 45两个输出分别对应0+1+2+3+4=10与0+1+...+9=45,验证了"数字流被完整聚合"的行为。
举一反三:更多的聚合模式
unstreamify方法内可以自由决定"如何消费流、返回什么",因此聚合逻辑远不止求和一种。结合源码 docstring 与 AWEL 的设计思路,常见的聚合模式包括:
- 统计个数:
async for逐条遍历并计数(前文MyUnstreamOperator即为此例); - 拼接文本:把流式返回的 token 片段拼成完整句子,这在流式 LLM 服务场景中非常实用——上游
StreamifyAbsOperator逐字产出,下游UnstreamifyAbsOperator拼装成完整回复; - 取最值/汇总:边遍历边维护
max、min、累加器等状态; - 构造容器:把流元素收集进
list、dict等数据结构后整体返回。
在教程序列中的位置与延伸阅读
本文对应 AWEL 教程"基础语法"部分第二节系列的第六篇。为形成完整的流式算子知识闭环,建议按以下顺序继续阅读:
- 前一篇 2.5_streamify_operator.md:讲解如何用
StreamifyAbsOperator把单数据展开成流,其中包含"模拟流式 LLM 服务"的完整示例,正是UnstreamifyAbsOperator最常见的上游搭档; - 后一篇 2.7_transform_stream_operator.md:讲解"流到流"变换,可与本文组合出"展开 → 逐条变换 → 聚合"的完整流水线;
- 深入源码可继续阅读 stream_operator.py(三类流算子的实现)与 task_impl.py(
SimpleTaskOutput/SimpleStreamTaskOutput的map/reduce/streamify/unstreamify/transform_stream五种变换能力); - 教程同目录下还有配套的 HTTP 触发(
network_program)与高级指南(advanced_guide)等章节,可进一步了解流式算子如何嵌入真实服务链路。
小结
UnstreamifyAbsOperator是 AWEL 流式编程中"收口"的关键算子:只需继承并覆写一个async def unstreamify(self, input_value: AsyncIterator[IN]) -> OUT方法,即可把上游数据流聚合为单个结果。其底层由_do_run驱动,经TaskOutput.unstreamify完成"流式输出 → 普通输出"的类型转换,配合StreamifyAbsOperator可以构建出完整的"展开—变换—聚合"异步流水线,是理解 DB-GPT 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),仅供参考