DB-GPT AWEL 实战使用 HttpTrigger 构建处理 POST 请求体的 HTTP 触发器【免费下载链接】DB-GPTopen-source agentic AI data assistant for the next generation of AI Data products.项目地址: https://gitcode.com/GitHub_Trending/db/DB-GPT本文基于 DB-GPTEosphoros-DB的 AWEL 工作流引擎教程讲解如何编写一个能接收并解析 JSON 请求体的 POST 型 HTTP 触发器。你将学会用 Pydantic 模型声明请求体结构、通过MapOperator对解析结果做转换并用setup_dev_environment一键起服务、用curl验证完整调用链同时结合dbgpt.core.awel的源码剖析请求体从 FastAPI 路由到 DAG 下游节点的完整流转机制。一、背景从 GET 到 POST在 AWEL 的网络化编程network program系列教程中3.1 基础 HTTP 触发器 演示了如何创建一个返回固定字符串的最简触发器3.2 处理 GET 请求 演示了如何根据 URL 查询参数返回结果本节3.3则处理POST 请求创建一个能根据请求体request body返回 JSON 响应的新 HTTP 触发器。与 GET 场景不同POST 的数据通常放在 JSON 请求体中因此HttpTrigger需要借助 Pydantic 模型来声明期望的请求体结构并在数据流入 DAG 后将其转换为模型实例供下游算子以属性访问如x.name的方式消费。后续若要处理流式响应可继续参考 3.4 流式响应。二、完整示例Say Hello To SomeonePOST 版在awel_tutorial目录下新建文件http_trigger_say_hello_post.py写入以下代码from dbgpt._private.pydantic import BaseModel, Field from dbgpt.core.awel import DAG, HttpTrigger, MapOperator, setup_dev_environment class TriggerReqBody(BaseModel): name: str Field(..., descriptionUser name) age: int Field(18, descriptionUser age) with DAG(awel_say_hello_post) as dag: trigger_task HttpTrigger( endpoint/awel_tutorial/say_hello_post, methodsPOST, request_bodyTriggerReqBody, status_code200 ) task MapOperator( map_functionlambda x: {message: fHello, {x.name}! You are {x.age} years old.} ) trigger_task task setup_dev_environment([dag], port5555)代码要点TriggerReqBody请求体模型用 Pydantic 声明请求体结构。name为必填字符串...表示 requiredage为整型、默认值18。注意这里从dbgpt._private.pydantic导入BaseModel与Field这是 DB-GPT 为兼容不同 Pydantic 大版本提供的私有代理模块能确保与 AWEL 内部的字段解析逻辑如field_is_required、model_fields一致。HttpTrigger声明 POST 端点endpoint/awel_tutorial/say_hello_post定义触发路径methodsPOST限定只接受 POSTrequest_bodyTriggerReqBody告诉触发器按该模型解析请求体status_code200为成功响应码。MapOperator做数据转换map_function接收触发器输出的数据这里返回一个 JSON 对象{message: ...}。下游拿到的是TriggerReqBody的实例因此可以直接x.name、x.age访问字段。建立依赖trigger_task task表示触发器输出流向 Map 算子。setup_dev_environment([dag], port5555)仅用于开发环境会自动创建 FastAPI 应用、注册触发器并用 uvicorn 在127.0.0.1:5555起服务详见下文源码剖析。执行代码poetry run python awel_tutorial/http_trigger_say_hello_post.py然后打开一个新终端向服务发送 POST 请求curl -X POST \ http://127.0.0.1:5555/api/v1/awel/trigger/awel_tutorial/say_hello_post \ -H Content-Type: application/json \ -d {name: John, age: 25}正常返回的 JSON 响应为{message:Hello, John! You are 25 years old.}验证完成后按CtrlC停止服务即可。如果请求体缺少必填字段namePydantic 校验会失败FastAPI 直接返回 422 校验错误不会进入 DAG 执行——这正是声明式请求体的价值入参在进流程之前就被结构化校验了。请求 URL 是怎么拼出来的注意最终请求地址是http://127.0.0.1:5555/api/v1/awel/trigger/awel_tutorial/say_hello_post而代码里只写了endpoint/awel_tutorial/say_hello_post。前缀/api/v1/awel/trigger来自HttpTriggerManager的默认router_prefix见 trigger_manager.pydef __init__( self, router: Optional[APIRouter] None, router_prefix: str /api/v1/awel/trigger, ) - None:注册触发器时register_trigger会用join_paths(self._router_prefix, real_endpoint)把前缀与端点拼成完整路径再挂载到 FastAPI 路由表上trigger_manager.py。所以你在curl中看到的 URL 实际上是router_prefix endpoint的组合。三、HttpTrigger 处理 POST 请求体的源码机制以下实现均位于 http_trigger.py帮助理解示例代码背后的调用链。3.1 构造函数参数校验与请求体归类HttpTrigger.__init__http_trigger.py的关键逻辑def __init__( self, endpoint: str, methods: Optional[Union[str, List[str]]] GET, request_body: Optional[RequestBody] None, http_trigger_body: Optional[Type[BaseHttpBody]] None, streaming_response: bool False, ... status_code: Optional[int] 200, ... ) - None: if not endpoint.startswith(/): endpoint / endpoint if not request_body and http_trigger_body: request_body http_trigger_body.get_body_class() streaming_predict_func http_trigger_body.streaming_predict_func() if not response_model and http_response_body: response_model http_response_body.get_body_class() ... self._methods [methods] if isinstance(methods, str) else methods可以看到endpoint若不以/开头会自动补前导斜杠methods传字符串会被包装成列表统一处理request_body支持多种形态见第五节。本例中request_body就是TriggerReqBodyPydantic 模型类。3.2 请求体到 DAG 数据的转换map方法POST 请求到达后FastAPI 已按request_body注解完成 JSON 反序列化。触发器自身的map方法http_trigger.py负责把原始输入规整为 Pydantic 实例async def map(self, input_data: Any) - Any: if not self._req_body or not input_data: return await super().map(input_data) if ( isinstance(self._req_body, type) and issubclass(self._req_body, BaseModel) and isinstance(input_data, dict) ): return self._req_body(**input_data) return await super().map(input_data)即当request_body是 Pydantic 模型子类、且输入还是dict时会执行self._req_body(**input_data)构造模型对象。这就是为什么下游MapOperator的 lambda 里能用x.name属性访问而不是x[name]。3.3 路由函数的动态生成_create_route_funchttp_trigger.py会区分“查询型方法”GET/DELETE参数走 query string与“体型方法”POST/PUT参数走 JSON body。对 POST 的分支非常简洁else: async def route_function(body: req_body_cls): # type: ignore return await _trigger_dag_func(body) route_function.__name__ name return route_function它直接生成一个签名带body: TriggerReqBody的 FastAPI 路由函数由 FastAPI 依据该类型注解完成请求体解析与校验随后调用_trigger_dag_func(body)进入 DAG。相比之下GET 场景需要遍历 Pydantic 字段、用inspect.Parameter动态重建函数签名把字段映射成 query 参数——这也是 GET/POST 两种触发器实现复杂度差异的根源。3.4 触发 DAG_trigger_dag与“单叶节点”约束_trigger_daghttp_trigger.py是真正执行流程的入口其中有一个关键约束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: with root_tracer.start_span(dbgpt.core.trigger.http.run_dag, span_id, metadatametadata): return await end_node.call(call_databody)也就是说HTTP 触发的 DAG必须有且只有一个叶子节点请求体会作为call_data从叶子节点开始反向驱动整条链路AWEL 的调用方向是“从输出端往回推”。非流式路径直接await end_node.call(call_databody)并把最终结果交给 FastAPI 序列化返回流式路径则走call_stream并包装成StreamingResponseSSE。执行过程还通过root_tracer生成 OpenTelemetry spandbgpt.core.trigger.http.run_dag便于链路追踪。3.5 MapOperator 在其中的角色下游的MapOperator定义于 common_operator.py其_do_run会把上游call_data包装后应用map_function再把输出写回 DAG 上下文。本例中它就是把TriggerReqBody实例映射为{message: ...}字典最终由 FastAPI 按 JSON 序列化输出——与示例的 curl 响应完全对应。3.6 开发环境启动链路setup_dev_environment定义在 dbgpt/core/awel/init.py其流程为配置日志默认写入dbgpt_awel_dev.log用_check_has_http_trigger(dags)检测是否存在HttpTrigger有则调用dbgpt.util.fastapi.create_app()创建 FastAPI 应用创建SystemApp并注册DefaultTriggerManager逐个遍历 DAG 的trigger_nodes调用trigger_manager.register_trigger(trigger, system_app)挂载路由最后uvicorn.run(app, hosthost, portport)启动服务默认127.0.0.1:5555。注意其 docstring 明确标注仅用于开发环境不适合生产环境生产场景应由 DB-GPT 的SystemApp组件体系统一管理触发器注册参考 initialize_awel。四、HttpTrigger 完整参数速查综合 http_trigger.py 的构造函数签名HttpTrigger支持以下参数参数类型默认值说明endpointstr必填API 端点路径不以/开头会自动补全支持{dag_id}占位符由_resolved_endpoint替换为实际 DAG IDmethodsstr \| List[str]GETHTTP 方法可传POST或[POST, PUT]等request_bodyRequest \| BaseModel \| Dict \| strNone请求体类型Pydantic 模型会被解析为结构化对象本文重点http_trigger_bodyType[BaseHttpBody]None通过内置 Body 资源类声明请求体如DictHttpBody、StringHttpBody、RequestHttpBody可同时携带流式判定函数streaming_responseboolFalse响应是否流式SSEstreaming_predict_funcCallableNone自定义“是否流式”的判定函数可依据请求体动态决定http_response_bodyType[BaseHttpBody]None响应体模型类会转换为 FastAPI 的response_modelresponse_modelTypeNone直接指定 FastAPI 响应模型response_headersDict[str, str]None额外响应头流式场景下默认注入text/event-stream等头response_media_typestrNone响应媒体类型status_codeint200成功状态码router_tagsList[str \| Enum]NoneOpenAPI 分组标签register_to_appboolFalseTrue时直接挂到 FastAPI app支持动态路由False时挂到触发器管理器的 Router 上其中trigger_modecommand或chat并非构造参数而是由_trigger_mode推断当request_body是CommonLLMHttpRequestBody子类时返回chat否则为commandhttp_trigger.py。五、request_body 的四种形态HttpTrigger的request_body联合类型定义于 http_trigger.pyRequestBody Union[Type[Request], Type[BaseModel], Type[Dict[str, Any]], Type[str]]结合 trigger/base.py 之外在 http_trigger.py 中注册的内置 Body 资源可以按需求选择BaseModel子类本文示例结构化、带校验适合参数明确的业务接口。GET/DELETE 场景下字段会变成 query 参数POST/PUT 场景下变成 JSON body。Dict[str, Any]请求体解析为字典适合字段不固定的透传接口。代码库还提供了开箱即用的DictHttpTriggerhttp_trigger.py它把request_bodydict、默认methodsPOST且register_to_appTrue固化好。str请求体解析为原始 JSON 字符串适合需要自行二次解析的场景对应便捷类StringHttpTriggerhttp_trigger.py。注意源码中明确限制GET/DELETE 查询型方法不支持str与dict类型会抛AWELHttpError。RequestStarlette 原始请求获取完整的 headers、query、原始 body自由度最高仓库中还提供了RequestHttpTrigger与面向 LLM 对话的CommonLLMHttpTrigger请求体为CommonLLMHttpRequestBody含model、messages、stream、conv_uid等字段并可将请求体转换为ModelRequestContext见 http_trigger.py。选择建议参数结构稳定且需要校验 →BaseModel字段不确定 →Dict[str, Any]要读 headers 或做协议级处理 →Request构建 LLM 对话网关 →CommonLLMHttpRequestBody。六、小结与延伸本文围绕 AWEL 教程 3.3 节完整演示了 POST 型 HTTP 触发器的编写与验证并结合源码梳理了调用链声明HttpTrigger(endpoint..., methodsPOST, request_bodyTriggerReqBody)定义端点、方法与 Pydantic 请求体挂载setup_dev_environment→DefaultTriggerManager→HttpTriggerManager以/api/v1/awel/trigger为前缀挂载路由并启动 uvicorn解析FastAPI 按request_body注解反序列化 JSONHttpTrigger.map保证下游拿到 Pydantic 实例执行_trigger_dag要求 DAG 单叶节点以call_databody驱动MapOperator等算子完成转换最后把叶子节点输出序列化为 JSON 返回约束与注意HttpTrigger仅支持单叶节点 DAGsetup_dev_environment仅限开发环境endpoint支持{dag_id}占位符以便多实例复用同一 DAG。基于此可以继续扩展为接口增加鉴权在路由函数前处理Request、用http_response_body约束响应结构、开启streaming_response实现 SSE 流式输出参考 3.4 流式响应或将多个 DAG 的触发器一起交给setup_dev_environment批量启动。【免费下载链接】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),仅供参考