DB-GPT AWEL 教程:InputOperator 输入算子——为 DAG 构建数据源头
【免费下载链接】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)工作流框架中的InputOperator(输入算子)展开。InputOperator 是 DAG 的起点节点,负责从用户自定义的输入源(InputSource)中读取数据并注入整条工作流。读完后,你将掌握如何使用内置的SimpleInputSource、SimpleCallDataInputSource快速搭建数据入口,理解单值数据与流式数据两种读取模式,并学会继承BaseInputSource编写自己的输入源。
什么是 InputOperator
InputOperator 是 AWEL 中的一种特殊算子:它没有任何输入,只拥有一个输出。按照 AWEL 基础语法教程 2.8 输入算子 的定位,它应当始终作为 DAG 中的第一个算子出现,其职责就是把某个输入源中的数据读取出来,交给下游算子消费。
从源码结构看,这一“无入参、单输出、读数据”的设计在实现上非常直白:InputOperator继承自BaseOperator,其核心运行逻辑_do_run只有三步——获取当前任务上下文、调用输入源的read方法读取数据、再通过可覆写的map钩子(默认原样返回)将数据写回上下文。相关实现位于 InputOperator 定义:
async def _do_run(self, dag_ctx: DAGContext) -> TaskOutput[OUT]: curr_task_ctx: TaskContext[OUT] = dag_ctx.current_task_context task_output = await self._input_source.read(curr_task_ctx) new_task_output: TaskOutput[OUT] = await task_output.map(self.map) curr_task_ctx.set_task_output(new_task_output) return task_output这意味着扩展输入能力的关键不在算子本身,而在输入源——只要实现一个符合InputSource接口的读取器,就能让 DAG 从任意数据渠道(常量、调用方传参、数据库、文件等)启动。
构建方式:把输入源交给 InputOperator
构造InputOperator的方式只有一种:将输入源实例传入构造函数。最小可用示例如下:
from dbgpt.core.awel import DAG, InputOperator, SimpleInputSource with DAG("awel_input_operator") as dag: input_source = SimpleInputSource(data="Hello, World!") input_task = InputOperator(input_source=input_source)这里SimpleInputSource(data="Hello, World!")创建了一个携带字符串数据的输入源,InputOperator(input_source=input_source)将其绑定为 DAG 的数据入口。
示例一:打印单值输入数据
第一个示例展示如何用InputOperator打印输入数据:使用以字符串为数据构建的SimpleInputSource,并通过MapOperator消费其输出。
新建文件awel_tutorial/input_operator_print_data.py,写入以下代码:
import asyncio from dbgpt.core.awel import DAG, MapOperator, InputOperator, SimpleInputSource with DAG("awel_input_operator") as dag: input_source = SimpleInputSource(data="Hello, World!") input_task = InputOperator(input_source=input_source) print_task = MapOperator(map_function=lambda x: print(x)) input_task >> print_task asyncio.run(print_task.call())执行:
poetry run python awel_tutorial/input_operator_print_data.py控制台输出:
Hello, World!几个要点:
input_task >> print_task用 AWEL 的>>语法声明了数据流向:输入算子先执行,其输出作为MapOperator的输入;print_task.call()触发整个 DAG 的异步执行,call是 BaseOperator 定义的异步入口,执行完毕后返回最终输出(本例中print副作用发生在图运行过程中)。
示例二:打印流式数据
第二个示例展示流式(streaming)数据场景:SimpleInputSource接收的是一个异步生成器,InputOperator会将其识别为流式输出,下游通过call_stream逐个消费元素。
新建文件awel_tutorial/input_operator_print_stream_data.py:
import asyncio from dbgpt.core.awel import DAG, InputOperator, SimpleInputSource async def stream_data(): for i in range(10): yield i with DAG("awel_input_operator") as dag: input_source = SimpleInputSource(data=stream_data()) input_task = InputOperator(input_source=input_source) async def print_stream(t: InputOperator): async for i in await t.call_stream(): print(i) asyncio.run(print_stream(input_task))执行:
poetry run python awel_tutorial/input_operator_print_stream_data.py输出为逐行打印的0到9:
0 1 2 3 4 5 6 7 8 9这个例子体现了 AWEL 流式编程能力的关键:call_stream(定义见 BaseOperator.call_stream)返回一个异步迭代器,数据边产生边消费,而不必一次性物化整个数据集。
从源码可以看到流式判断的实际机制:BaseInputSource.read在读取数据后,若构造时未显式指定streaming参数,就会自动探测数据是否为(异步)迭代器——是则包装为SimpleStreamTaskOutput,否则包装为普通SimpleTaskOutput;若显式指定了streaming,则以显式值为准。完整逻辑见 BaseInputSource.read:
data = self._read_data(task_ctx) if self._streaming_data is None: streaming_data = _is_async_iterator(data) or _is_iterator(data) else: streaming_data = self._streaming_data if streaming_data: if self._is_read: raise ValueError(f"Input iterator {data} has been read!") it_data = _to_async_iterator(data) output: TaskOutput = SimpleStreamTaskOutput(it_data) else: output = SimpleTaskOutput(data) self._is_read = True return output这里有一个值得注意的约束:流式输入源只能被读取一次。若同一个流式输入源所在的任务被再次执行,read会抛出ValueError("Input iterator ... has been read!")。这与迭代器“消费即耗尽”的语义一致,在复用 DAG 时应保持这一预期。
示例三:打印 Call Data
Call data指调用算子call或call_stream方法时传入的数据。该示例使用SimpleCallDataInputSource——一个不携带固定数据、而是在运行时从调用参数中取数的输入源,因此同一个 DAG 可以在不同次调用中注入不同数据。
新建文件awel_tutorial/input_operator_print_call_data.py:
import asyncio from dbgpt.core.awel import DAG, MapOperator, InputOperator, SimpleCallDataInputSource with DAG("awel_input_operator") as dag: input_source = SimpleCallDataInputSource() input_task = InputOperator(input_source=input_source) print_task = MapOperator(map_function=lambda x: print(x)) input_task >> print_task asyncio.run(print_task.call(call_data="Hello, World!")) asyncio.run(print_task.call(call_data="AWEL is cool!"))执行:
poetry run python awel_tutorial/input_operator_print_call_data.py输出:
Hello, World! AWEL is cool!可以看到同一份 DAG 定义被运行了两次,每次call(call_data=...)传入的数据都成为了 DAG 的起始数据。
从实现层面看,SimpleCallDataInputSource的_read_data从任务上下文中取出call_data字典,并读取其中的"data"键;如果取不到数据,会抛出ValueError("No call data for current SimpleCallDataInputSource")。相关代码见 SimpleCallDataInputSource。而任务上下文的call_data属性则是通过TaskContext的元数据(metadata)存取,见 TaskContext.call_data。可以推断:call()调用时传入的call_data会被注入到执行上下文的元数据中,供任何需要它的算子(包括输入源)读取,这也是 AWEL 中“运行时参数化”DAG 的基础机制。
内置输入源小结
AWEL 提供两种内置输入源:
| 输入源 | 用途 |
|---|---|
SimpleInputSource | 用单个数据或流式数据创建输入源,数据在构造时确定(data参数) |
SimpleCallDataInputSource | 数据来自算子call/call_stream方法传入的call data,运行时确定 |
两者都实现自抽象基类InputSource(见 InputSource 定义),该基类除抽象方法read外,还提供三个便捷的类方法用于快速创建输入源:
InputSource.from_data(data):以单个数据创建,内部等价于SimpleInputSource(data, streaming=False);InputSource.from_iterable(iterable):以可迭代对象创建,内部等价于SimpleInputSource(iterable, streaming=True);InputSource.from_callable():以调用数据创建,内部等价于SimpleCallDataInputSource()。
值得一提的是,InputOperator还提供了一个类方法dummy_input,用于创建一个携带“占位数据”(默认为SKIP_DATA)的假输入算子,便于在图中占位而不引入真实数据,见 InputOperator.dummy_input。从源码结构看,项目中的TriggerOperator(DAG 触发器)也是直接继承InputOperator并内置SimpleCallDataInputSource实现的,可见输入算子体系在 AWEL 触发机制中的基础性地位。
创建自定义输入源
创建自己的输入源最简单的方式是继承BaseInputSource并重写_read_data方法。下面的示例实现了一个返回固定字符串的输入源:
import asyncio from dbgpt.core.awel import DAG, InputOperator, MapOperator, BaseInputSource, TaskContext class MyInputSource(BaseInputSource): """Create an input source with a single data""" def _read_data(self, ctx: TaskContext) -> str: return "Hello, World!" with DAG("awel_input_operator") as dag: input_source = MyInputSource() input_task = InputOperator(input_source=input_source) print_task = MapOperator(map_function=lambda x: print(x)) input_task >> print_task asyncio.run(print_task.call())将文件保存为awel_tutorial/my_input_source.py后执行:
poetry run python awel_tutorial/my_input_source.py输出:
Hello, World!_read_data(self, ctx: TaskContext)接收一个TaskContext参数,这为自定义输入源留下了充足的扩展空间:你完全可以在其中读取运行时上下文、查询外部系统或发起网络请求。例如返回一个生成器即可得到一个自定义流式输入源——因为如前所述,BaseInputSource.read会对返回值做迭代器自动探测。BaseInputSource的完整定义(含streaming构造参数与read的默认实现)见 BaseInputSource。
小结:InputOperator 在 AWEL DAG 中的位置
结合本教程与源码,可以把 InputOperator 的要点归纳为:
- DAG 的起点:InputOperator 无输入、单输出,习惯上作为 DAG 第一个算子,负责把外部数据“接”进工作流;
- 数据形态二选一:单值数据(
SimpleTaskOutput)与流式数据(SimpleStreamTaskOutput),流式源只能读取一次; - 参数化运行:通过
SimpleCallDataInputSource+call(call_data=...),同一 DAG 定义可在多次运行中携带不同输入,是 AWEL 实现运行时参数化的关键手段; - 面向扩展设计:只需继承
BaseInputSource重写_read_data,就能把工作流的输入端对接到任意自定义数据渠道。
如果你正在把业务逻辑接入 AWEL,建议从SimpleInputSource起步验证链路,再用SimpleCallDataInputSource完成运行时参数化,最终按需实现自定义输入源对接真实数据系统。更多基础语法(如 MapOperator、ReduceOperator、JoinOperator 等)可参考同目录下的 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),仅供参考