DB-GPT AWEL 教程:用 UnstreamifyAbsOperator 将数据流聚合为单个结果
2026/9/14 19:00:52 网站建设 项目流程

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]并带有两个泛型参数INOUT(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 还提供了通用的MapOperatorReduceOperator等算子(参见 2.2_reduce_operator.md)。UnstreamifyAbsOperator与它们的本质区别在于输入形态:普通算子拿到的是父节点输出的一条具体数据,而UnstreamifyAbsOperator的输入是AsyncIterator[IN],即必须通过async for逐条消费的异步迭代器。因此它天然适合接在StreamifyAbsOperatorTransformStreamAbsOperator等流式节点之后,作为整个流式链路的"终点聚合器"。

自定义一个 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()

这段代码做了三件事:

  1. 通过from dbgpt.core.awel import ...引入DAGUnstreamifyAbsOperator。这两个符号在 packages/dbgpt-core/src/dbgpt/core/awel/__init__.py 中均被显式导出(UnstreamifyAbsOperator同时出现在__all__列表中)。
  2. 定义SumOperator,泛型参数[int, int]表示"输入数据流元素为 int、聚合结果为 int"。
  3. 覆写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

这条链路包含三个关键环节:

  1. 取上游输出:通过curr_task_ctx.task_input.parent_outputs[0].task_output拿到父节点的TaskOutput。注意它固定取第一个父输出,因此UnstreamifyAbsOperator是单输入算子,DAG 中只能有一个上游节点直接连接到它;
  2. 调用 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方法,再把返回值包装成一个非流式的SimpleTaskOutputUnStreamFunc的类型别名定义在 packages/dbgpt-core/src/dbgpt/core/awel/task/base.py:

UnStreamFunc = Callable[[AsyncIterator[IN]], OUT]
  1. 写回任务上下文:聚合出的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 等章节对连线语法有更系统的介绍;
  • 调用终点是SumOperatorsum_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=100+1+...+9=45,验证了"数字流被完整聚合"的行为。

举一反三:更多的聚合模式

unstreamify方法内可以自由决定"如何消费流、返回什么",因此聚合逻辑远不止求和一种。结合源码 docstring 与 AWEL 的设计思路,常见的聚合模式包括:

  • 统计个数async for逐条遍历并计数(前文MyUnstreamOperator即为此例);
  • 拼接文本:把流式返回的 token 片段拼成完整句子,这在流式 LLM 服务场景中非常实用——上游StreamifyAbsOperator逐字产出,下游UnstreamifyAbsOperator拼装成完整回复;
  • 取最值/汇总:边遍历边维护maxmin、累加器等状态;
  • 构造容器:把流元素收集进listdict等数据结构后整体返回。

在教程序列中的位置与延伸阅读

本文对应 AWEL 教程"基础语法"部分第二节系列的第六篇。为形成完整的流式算子知识闭环,建议按以下顺序继续阅读:

  • 前一篇 2.5_streamify_operator.md:讲解如何用StreamifyAbsOperator把单数据展开成流,其中包含"模拟流式 LLM 服务"的完整示例,正是UnstreamifyAbsOperator最常见的上游搭档;
  • 后一篇 2.7_transform_stream_operator.md:讲解"流到流"变换,可与本文组合出"展开 → 逐条变换 → 聚合"的完整流水线;
  • 深入源码可继续阅读 stream_operator.py(三类流算子的实现)与 task_impl.py(SimpleTaskOutput/SimpleStreamTaskOutputmap/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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询