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 请求体驱动的流式(Streaming)HTTP 接口:从编写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(..., description="The 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", methods="POST", request_body=TriggerReqBody, status_code=200, streaming_predict_func=lambda x: True ) task = NumberProducerOperator() trigger_task >> task setup_dev_environment([dag], port=5555)代码中有四个关键点:
TriggerReqBody:Pydantic 模型定义了 POST 请求体的结构,这里只含一个必填字段n(要流式输出的整数个数)。HttpTrigger会用它做请求体校验与反序列化;NumberProducerOperator:继承StreamifyAbsOperator[TriggerReqBody, int],实现抽象方法streamify,把输入值转换为一个AsyncIterator[int],逐个yield数字(教程中每个数字后拼接\n,因此 curl 输出每行一个数字);HttpTrigger:endpoint指定接口路径,methods="POST"指定方法,request_body=TriggerReqBody指定请求体模型,status_code=200是响应状态码,streaming_predict_func=lambda x: True是一个流式判定函数——它始终返回True,即无论请求内容如何,本接口永远走流式响应;setup_dev_environment([dag], port=5555):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完成后按Ctrl+C停止服务即可。
注意最终 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__的完整签名如下:
| 参数 | 类型 | 默认值 | 说明 |
|---|---|---|---|
endpoint | str | 必填 | 接口路径,不以/开头会自动补全 |
methods | str \| List[str] | "GET" | 允许的 HTTP 方法,支持GET/PUT/POST/DELETE |
request_body | Type | None | 请求体模型,可为 Pydantic 模型、dict、str或starlette.Request |
http_trigger_body | Type[BaseHttpBody] | None | 使用内置 Body 封装类(见下文),自动派生request_body与流式判定函数 |
streaming_response | bool | False | 是否强制流式响应 |
streaming_predict_func | Callable[[CommonRequestType], bool] | None | 动态流式判定函数,按请求内容决定是否流式 |
http_response_body | Type[BaseHttpBody] | None | 响应体模型,自动派生response_model |
response_model | Type | None | FastAPI 响应模型 |
response_headers | Dict[str, str] | None | 自定义响应头(流式时生效) |
response_media_type | str | None | 响应媒体类型,流式时默认text/event-stream |
status_code | int | 200 | 响应状态码 |
router_tags | List[str] | None | OpenAPI 分组标签 |
register_to_app | bool | False | True时直接挂到 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_func:
def _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], port=5555)的完整实现见 awel/init.py。其工作流程为:
- 调用
setup_logging初始化日志(默认写dbgpt_awel_dev.log); - 通过
_check_has_http_trigger检测 DAG 中是否存在HttpTrigger,存在则用create_app()创建 FastAPI 应用; - 创建
SystemApp与DefaultTriggerManager,对每个 DAG:默认show_dag_graph=True会调用dag.visualize_dag()生成可视化图(依赖 graphviz,缺失时仅告警,不影响运行),然后逐个把dag.trigger_nodes注册进触发器管理器; - 若管理器判定需要持续运行(
keep_running(),即存在已注册的 HTTP 触发器),最后以uvicorn.run(app, host, port)启动服务——这就是为什么 Ctrl+C 能停掉整个服务器。
注意文档字符串明确说明:该函数仅用于开发环境,不适用于生产环境("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 已注册过会抛ValueError("Route {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_app=True的内置触发器(如后文的DictHttpTrigger)则直接mount_to_app,以priority=10动态加入 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_data=body) 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_data=body) ... return StreamingResponse( trace_generator, headers=headers, media_type=media_type, background=background_tasks, )从源码可以确认四个关键实现事实:
- 单叶子节点约束:HttpTrigger 触发的 DAG 有且只能有一个叶子节点,流式数据源就是该节点——教程示例的
trigger_task >> task结构正好满足; - 非流式:直接
await end_node.call(call_data=body),一次性返回终值; - 流式:调用
end_node.call_stream(call_data=body)拿到异步迭代器,包装成 FastAPIStreamingResponse,默认响应头为text/event-stream+no-cache+keep-alive+chunked,media_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_call=True执行 workflow,再检查任务输出task_output.is_stream——是流则直接取output_stream迭代器,否则把单一输出包装成单次 yield 的生成器。这就解释了为什么NumberProducerOperator用streamify逐块 yield 的字符串最终能变成 HTTP 响应体上一行一行的0..4:每个yield的块被原样转发,直到生成器耗尽,背景任务收尾、连接关闭。
五、内置触发器变体与其他流式算子
除了教程使用的通用HttpTrigger,同一模块还内置了几个便捷子类,适合不同场景:
DictHttpTrigger(第 823 行起):请求体按dict解析,方法默认POST,内部强制register_to_app=True;StringHttpTrigger:请求体按 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 使用相同
endpoint+methods会在注册阶段直接抛错,启动日志会提示 "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),仅供参考