ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

DB-GPT AWEL 教程:InputOperator 输入算子——为 DAG 构建数据源头

DB-GPT AWEL 教程:InputOperator 输入算子——为 DAG 构建数据源头 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 的 AWELAgentic Workflow Expression Language工作流框架中的InputOperator输入算子展开。InputOperator 是 DAG 的起点节点负责从用户自定义的输入源InputSource中读取数据并注入整条工作流。读完后你将掌握如何使用内置的SimpleInputSource、SimpleCallDataInputSource快速搭建数据入口理解单值数据与流式数据两种读取模式并学会继承BaseInputSource编写自己的输入源。什么是 InputOperatorInputOperator 是 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(dataHello, World!) input_task InputOperator(input_sourceinput_source)这里SimpleInputSource(dataHello, World!)创建了一个携带字符串数据的输入源InputOperator(input_sourceinput_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(dataHello, World!) input_task InputOperator(input_sourceinput_source) print_task MapOperator(map_functionlambda 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.pyimport 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(datastream_data()) input_task InputOperator(input_sourceinput_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到90 1 2 3 4 5 6 7 8 9这个例子体现了 AWEL 流式编程能力的关键call_stream定义见 BaseOperator.call_stream返回一个异步迭代器数据边产生边消费而不必一次性物化整个数据集。从源码可以看到流式判断的实际机制BaseInputSource.read在读取数据后若构造时未显式指定streaming参数就会自动探测数据是否为异步迭代器——是则包装为SimpleStreamTaskOutput否则包装为普通SimpleTaskOutput若显式指定了streaming则以显式值为准。完整逻辑见 BaseInputSource.readdata 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(fInput 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 DataCall data指调用算子call或call_stream方法时传入的数据。该示例使用SimpleCallDataInputSource——一个不携带固定数据、而是在运行时从调用参数中取数的输入源因此同一个 DAG 可以在不同次调用中注入不同数据。新建文件awel_tutorial/input_operator_print_call_data.pyimport asyncio from dbgpt.core.awel import DAG, MapOperator, InputOperator, SimpleCallDataInputSource with DAG(awel_input_operator) as dag: input_source SimpleCallDataInputSource() input_task InputOperator(input_sourceinput_source) print_task MapOperator(map_functionlambda x: print(x)) input_task print_task asyncio.run(print_task.call(call_dataHello, World!)) asyncio.run(print_task.call(call_dataAWEL 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, streamingFalse)InputSource.from_iterable(iterable)以可迭代对象创建内部等价于SimpleInputSource(iterable, streamingTrue)InputSource.from_callable()以调用数据创建内部等价于SimpleCallDataInputSource()。值得一提的是InputOperator还提供了一个类方法dummy_input用于创建一个携带“占位数据”默认为SKIP_DATA的假输入算子便于在图中占位而不引入真实数据见 InputOperator.dummy_input。从源码结构看项目中的TriggerOperatorDAG 触发器也是直接继承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_sourceinput_source) print_task MapOperator(map_functionlambda 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流式源只能读取一次参数化运行通过SimpleCallDataInputSourcecall(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),仅供参考
返回列表