DB-GPT AWEL 流式 HTTP 触发器实战:基于 HttpTrigger 与 StreamifyAbsOperator 构建 SSE 流式接口

📅 发布时间:2026/9/13 18:35:57
DB-GPT AWEL 流式 HTTP 触发器实战:基于 HttpTrigger 与 StreamifyAbsOperator 构建 SSE 流式接口
DB-GPT AWEL 流式 HTTP 触发器实战基于 HttpTrigger 与 StreamifyAbsOperator 构建 SSE 流式接口【免费下载链接】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 官方教程 3.4 节讲解如何用HttpTrigger搭建一个基于 POST 请求体驱动的流式StreamingHTTP 接口从编写StreamifyAbsOperator算子、配置streaming_predict_func到本地起服务、用 curl 验证逐行输出的数字流并结合 http_trigger.py 源码解析路由注册、流式判定与 SSE 响应的完整链路。读完你可以独立完成一个生产可用的 AWEL 流式 API并理解其底层实现细节。一、实战示例流式输出 0 到 n-1 的数字教程的原始目标是创建一个返回流式响应的 HTTP 触发器根据 POST 请求体决定是否流式输出。在awel_tutorial目录下新建文件http_trigger_stream_numbers.py完整代码如下即教程原始示例可直接复制运行from dbgpt._private.pydantic import BaseModel, Field from dbgpt.core.awel import DAG, HttpTrigger, StreamifyAbsOperator, setup_dev_environment from typing import AsyncIterator class TriggerReqBody(BaseModel): n: int Field(..., descriptionThe number of integers to be streamed) class NumberProducerOperator(StreamifyAbsOperator[TriggerReqBody, int]): Create a stream of numbers from 0 to n-1 async def streamify(self, req: TriggerReqBody) - AsyncIterator[int]: for i in range(req.n): yield str(i) \n with DAG(awel_stream_numbers) as dag: trigger_task HttpTrigger( endpoint/awel_tutorial/stream_numbers, methodsPOST, request_bodyTriggerReqBody, status_code200, streaming_predict_funclambda x: True ) task NumberProducerOperator() trigger_task task setup_dev_environment([dag], port5555)代码中有四个关键点TriggerReqBodyPydantic 模型定义了 POST 请求体的结构这里只含一个必填字段n要流式输出的整数个数。HttpTrigger会用它做请求体校验与反序列化NumberProducerOperator继承StreamifyAbsOperator[TriggerReqBody, int]实现抽象方法streamify把输入值转换为一个AsyncIterator[int]逐个yield数字教程中每个数字后拼接\n因此 curl 输出每行一个数字HttpTriggerendpoint指定接口路径methodsPOST指定方法request_bodyTriggerReqBody指定请求体模型status_code200是响应状态码streaming_predict_funclambda x: True是一个流式判定函数——它始终返回True即无论请求内容如何本接口永远走流式响应setup_dev_environment([dag], port5555)AWEL 提供的开发环境启动器在127.0.0.1:5555启动一个 FastAPI uvicorn 服务并注册该 DAG 的所有触发器。运行代码poetry run python awel_tutorial/http_trigger_stream_numbers.py然后另开一个终端向服务发送 POST 请求curl -X POST \ http://127.0.0.1:5555/api/v1/awel/trigger/awel_tutorial/stream_numbers \ -H Content-Type: application/json \ -d {n: 5}预期输出0 1 2 3 4完成后按CtrlC停止服务即可。注意最终 URL 的构成/api/v1/awel/trigger是 AWEL 触发器管理器的固定路由前缀/awel_tutorial/stream_numbers才是你在HttpTrigger(endpoint...)里声明的路径两者拼接后才是真实可访问的接口。这一点在源码 trigger_manager.py 中可以确认HttpTriggerManager的默认参数就是router_prefix/api/v1/awel/trigger注册触发器时会用join_paths(self._router_prefix, real_endpoint)拼接完整路径。仓库中 simple_rag_summary_example.py 等其他 AWEL 示例也使用了同样的 URL 模式可以佐证该约定是全局一致的。二、HttpTrigger 参数详解源码级教程示例只用了HttpTrigger的一小部分参数。http_trigger.py 中HttpTrigger.__init__的完整签名如下参数类型默认值说明endpointstr必填接口路径不以/开头会自动补全methodsstr \| List[str]GET允许的 HTTP 方法支持GET/PUT/POST/DELETErequest_bodyTypeNone请求体模型可为 Pydantic 模型、dict、str或starlette.Requesthttp_trigger_bodyType[BaseHttpBody]None使用内置 Body 封装类见下文自动派生request_body与流式判定函数streaming_responseboolFalse是否强制流式响应streaming_predict_funcCallable[[CommonRequestType], bool]None动态流式判定函数按请求内容决定是否流式http_response_bodyType[BaseHttpBody]None响应体模型自动派生response_modelresponse_modelTypeNoneFastAPI 响应模型response_headersDict[str, str]None自定义响应头流式时生效response_media_typestrNone响应媒体类型流式时默认text/event-streamstatus_codeint200响应状态码router_tagsList[str]NoneOpenAPI 分组标签register_to_appboolFalseTrue时直接挂到 FastAPI app 上支持动态路由两个与流式直接相关的参数需要重点区分streaming_response静态开关。设为True时该接口的所有请求都走流式分支streaming_predict_func动态开关。接收本次请求的 body教程中是TriggerReqBody实例返回bool决定这一次请求是否流式。教程示例中传入lambda x: True因此恒为流式。判定优先级见路由函数_trigger_dag_func先取self._streaming_response作为基础值若配置了streaming_predict_func则用其返回值覆盖再否则若 body 是BaseHttpBody实例则回落到默认判定函数。默认流式判定函数读取 body 中的 streaming 字段不传streaming_predict_func时源码会走_default_streaming_predict_funcdef _default_streaming_predict_func(body: CommonRequestType) - bool: if isinstance(body, BaseModel): body model_to_dict(body) elif isinstance(body, str): try: body json.loads(body) except Exception: return False elif not isinstance(body, dict): return False streaming body.get(streaming) or body.get(stream) return _parse_bool(streaming)也就是说把 body 归一化为 dict 后读取streaming或stream字段并按布尔解析。这实际上是一个按请求协商流式的机制——同一个接口既支持一次性返回也支持流式返回由客户端在请求体里声明。若你在TriggerReqBody中自行加了stream: bool False字段去掉streaming_predict_func后客户端传{n: 5, stream: true}即可获得流式输出。GET 请求体与 POST 请求体的处理差异_create_route_func中有一个is_query_method分支当方法全部为GET/DELETE时Pydantic 请求体模型会被展开为查询参数每个 field 映射为一个 query 参数并重建函数签名供 FastAPI/OpenAPI 使用而POST/PUT则直接把request_body_cls作为 JSON body 类型。dict或str类型的请求体不支持 query 方法源码中会直接抛出AWELHttpError。触发器拿到 body 后还会经过HttpTrigger.map第 565-584 行做最终转换若请求体是BaseModel子类且输入是 dict会实例化为对应的 Pydantic 对象再传入 DAG——这就是streamify能直接收到TriggerReqBody实例的原因。三、StreamifyAbsOperator把值变成流教程中的NumberProducerOperator继承自 StreamifyAbsOperator它的定义是“将一个IN值转换为AsyncIterator[OUT]”的抽象算子class StreamifyAbsOperator(BaseOperator[OUT], ABC, Generic[IN, OUT]): An abstract operator that converts a value of IN to an AsyncIterator[OUT]. streaming_operator True async def _do_run(self, dag_ctx: DAGContext) - TaskOutput[OUT]: ... output await wrapped_call_data.streamify(self.streamify) curr_task_ctx.set_task_output(output) return output abstractmethod async def streamify(self, input_value: IN) - AsyncIterator[OUT]: ...实现细节必须实现抽象方法streamify(input_value: IN) - AsyncIterator[OUT]内部用async生成器yield每个数据块类属性streaming_operator True声明这是一个流式算子其任务输出会被包装成流式TaskOutput_do_run中通过wrapped_call_data.streamify(self.streamify)完成同一模块还提供了两个配套抽象算子可组合出完整的流式处理链UnstreamifyAbsOperator把AsyncIterator[IN]聚合回单个OUT值例如统计流中元素个数TransformStreamAbsOperator流到流的逐元素转换例如对每个元素做1。在本例的 DAG 中NumberProducerOperator是唯一叶子节点它的streamify输出的迭代器就是最终 HTTP 响应的数据源。四、服务端全链路从请求到 SSE 响应1. setup_dev_environment开发环境如何起服务教程最后一行setup_dev_environment([dag], port5555)的完整实现见 awel/init.py。其工作流程为调用setup_logging初始化日志默认写dbgpt_awel_dev.log通过_check_has_http_trigger检测 DAG 中是否存在HttpTrigger存在则用create_app()创建 FastAPI 应用创建SystemApp与DefaultTriggerManager对每个 DAG默认show_dag_graphTrue会调用dag.visualize_dag()生成可视化图依赖 graphviz缺失时仅告警不影响运行然后逐个把dag.trigger_nodes注册进触发器管理器若管理器判定需要持续运行keep_running()即存在已注册的 HTTP 触发器最后以uvicorn.run(app, host, port)启动服务——这就是为什么 CtrlC 能停掉整个服务器。注意文档字符串明确说明该函数仅用于开发环境不适用于生产环境Just using in development environment, not production environment。生产场景应把触发器注册到 DB-GPT 应用自身的SystemApp通过initialize_awel走DAGManager加载 DAG 目录。2. 路由注册/api/v1/awel/trigger 前缀从哪里来HttpTriggerManager 负责把所有HttpTrigger挂到统一前缀下register_trigger中先取trigger._resolved_endpoint()支持{dag_id}占位符替换为真实 DAG ID再join_paths(self._router_prefix, real_endpoint)拼出完整路径并做重复路由检测——同一 path 下相同 method 已注册过会抛ValueErrorRoute {path} method {m} already registered这也是启动服务前要先确认端口未被占用的原因之一trigger.register_to_app()返回False时教程示例即此情况走mount_to_router路由先挂到APIRouter待所有触发器注册完后在_init_app里以app.include_router(router, prefix/api/v1/awel/trigger, tags[AWEL])一次性挂载register_to_appTrue的内置触发器如后文的DictHttpTrigger则直接mount_to_app以priority10动态加入 app并重置openapi_schema与middleware_stack缓存。3. 流式响应的真正出口_trigger_dag请求到达后动态路由函数最终调用_trigger_dag这里集中体现了流式/非流式两条路径的差异leaf_nodes dag.leaf_nodes if len(leaf_nodes) ! 1: raise ValueError(HttpTrigger just support one leaf node in dag) end_node cast(BaseOperator, leaf_nodes[0]) ... if not streaming_response: ... return await end_node.call(call_databody) else: headers response_headers media_type response_media_type if response_media_type else text/event-stream if not headers: headers { Content-Type: text/event-stream, Cache-Control: no-cache, Connection: keep-alive, Transfer-Encoding: chunked, } _generator await end_node.call_stream(call_databody) ... return StreamingResponse( trace_generator, headersheaders, media_typemedia_type, backgroundbackground_tasks, )从源码可以确认四个关键实现事实单叶子节点约束HttpTrigger 触发的 DAG 有且只能有一个叶子节点流式数据源就是该节点——教程示例的trigger_task task结构正好满足非流式直接await end_node.call(call_databody)一次性返回终值流式调用end_node.call_stream(call_databody)拿到异步迭代器包装成 FastAPIStreamingResponse默认响应头为text/event-streamno-cachekeep-alivechunkedmedia_type默认text/event-stream可通过response_media_type参数覆盖收尾与追踪DAG 结束回调dag._after_dag_end被放入BackgroundTasks在流耗尽后执行资源清理整个流经过root_tracer.wrapper_async_stream包装流式输出同样进入 OpenTelemetry 式 span 追踪span 名dbgpt.core.trigger.http.run_dag。4. call_stream算子层的流式执行入口叶子节点的call_stream是这条链路的核心它把call_data包成{data: ...}后以streaming_callTrue执行 workflow再检查任务输出task_output.is_stream——是流则直接取output_stream迭代器否则把单一输出包装成单次 yield 的生成器。这就解释了为什么NumberProducerOperator用streamify逐块 yield 的字符串最终能变成 HTTP 响应体上一行一行的0..4每个yield的块被原样转发直到生成器耗尽背景任务收尾、连接关闭。五、内置触发器变体与其他流式算子除了教程使用的通用HttpTrigger同一模块还内置了几个便捷子类适合不同场景DictHttpTrigger第 823 行起请求体按dict解析方法默认POST内部强制register_to_appTrueStringHttpTrigger请求体按 JSON 字符串解析同样默认POST 动态挂载CommonLLMHttpTrigger配合CommonLLMHttpRequestBody含model、messages、stream、temperature等字段使用的 LLM 通用触发器触发模式会被识别为chat见_trigger_mode并额外输出request_string_messages等映射字段。如果需要在流式链路末端做聚合如把整段 LLM 流拼成完整文本后再存储可接一个UnstreamifyAbsOperator若需要逐块改写翻译、脱敏、加前缀用TransformStreamAbsOperator。三者组合即为 AWEL 流式 DAG 的完整算子族。六、验证与排错要点URL 拼接忘记/api/v1/awel/trigger前缀是最常见的 404 原因完整路径 前缀 endpoint可用{dag_id}占位符路由冲突同一进程内两个 DAG 使用相同endpointmethods会在注册阶段直接抛错启动日志会提示 Route ... already registered多叶子节点DAG 末端若分叉出多个无出边节点_trigger_dag会抛ValueError: HttpTrigger just support one leaf node in dag流式判定若希望同一接口按请求决定流式与否删除streaming_predict_func改用 body 中的streaming/stream字段默认判定函数的行为或自己传入判定函数仅开发用途setup_dev_environment的 docstring 明确限定开发环境生产部署应复用 DB-GPT 应用的SystemApp与DAGManager机制graphviz 告警启动时若看到 DAG 可视化失败告警属于show_dag_graph缺 graphviz 的正常提示可pip install graphviz解决不影响接口功能。小结教程 3.4 节用一个 30 行不到的示例完整覆盖了 AWEL 流式 HTTP 接口的三要素请求体模型Pydantic→ 流式算子StreamifyAbsOperator.streamify→ 流式判定streaming_predict_func / streaming_response。结合源码可以看到HttpTrigger本质上是一个把 FastAPI 路由动态织入 DAG 生命周期的入口算子而真正的流式出口统一收敛在_trigger_dag的StreamingResponse分支中。掌握这套机制后你可以把示例中的NumberProducerOperator替换为任意异步生成逻辑LLM token 流、检索结果流、数据批处理进度等以同样的方式发布为 DB-GPT 生态中的流式 API。【免费下载链接】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),仅供参考