ARTICLE DETAIL

建站实战干货

来自一线的建站与推广经验沉淀,每一条都经过真实交付验证。

Apache Airflow Common AI Provider 架构解析:基于 pydantic-ai 的 LLM 接入设计、Toolset 扩展规范与安全边界

2026/9/14 20:51:53 拓冰建站 浏览量
Apache Airflow Common AI Provider 架构解析:基于 pydantic-ai 的 LLM 接入设计、Toolset 扩展规范与安全边界 Apache Airflow Common AI Provider 架构解析基于 pydantic-ai 的 LLM 接入设计、Toolset 扩展规范与安全边界【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow本文以providers/common/ai/AGENTS.md中定义的 Common AI Provider 开发指令为主体逐条展开其设计原则、新增 LLM Provider 与 Toolset 的标准流程、安全红线与常见陷阱并结合PydanticAIHook、SQLToolset、SQL 校验工具与AgentOperator的源码实现说明这些约定在 Apache Airflow 中是如何落地、被验证和执行的。一、Common AI Provider 的定位一条薄桥Common AI Provider发布包名apache-airflow-providers-common-ai见 provider.yaml的职责是把 Airflow 流水线接入 LLM。根据 AGENTS.md 的定义它通过封装 pydantic-ai 来连接各类大模型服务而连接层Hook的定位极其明确The hook is a thin bridge between Airflow connections and pydantic-ais model/provider abstractions. Hook 是 Airflow 连接与 pydantic-ai 模型/Provider 抽象之间的一条薄桥。pydantic-ai 自身已经支持 20 家模型供应商OpenAI、Anthropic、Google、Azure、Bedrock、Ollama 等通过infer_model()和AzureProvider、BedrockProvider等 Provider 类完成模型解析。这一分工在源码中可以直接看到hooks/pydantic_ai.py 顶部只从 pydantic-ai 导入Agent、infer_model、infer_provider、infer_provider_class自身没有任何对 OpenAI/Anthropic SDK 的直接依赖。这种薄桥设计的收益是Airflow 侧不为任何一家模型厂商写专属逻辑供应商适配的成本、升级和 bug 都由上游 pydantic-ai 承担。AGENTS.md 因此开宗明义地给出第一条设计原则——Delegate to pydantic-ai凡是 pydantic-ai 已经处理的供应商特定逻辑一律不要在本包中重新实现写新代码前先核对 pydantic-ai 的模型支持清单。二、五条核心设计原则逐条对照源码AGENTS.md 的Design Principles一节是整个 provider 的架构宪法共五条。以下逐条说明并给出仓库内的实现证据。2.1 委托给 pydantic-ai不做重复实现Hook 的模型解析完全委托给infer_model()。在 PydanticAIHook.get_conn() 中可以看到两级解析顺序显式凭据当_get_provider_kwargs()返回非空 dict 时用infer_provider_class(pname)(**kwargs)实例化 Provider 并包进provider_factory再交给infer_model(model_name, provider_factory...)若 Provider 类拒绝这些 kwargsTypeError会记录警告并回退到环境变量认证infer_provider(pname)默认解析连接中没有可用凭据时直接infer_model(model_name)由 pydantic-ai 读取OPENAI_API_KEY、AWS_PROFILE等标准环境变量。解析出的Model会被缓存在 hook 实例中self._model同一实例内不会重复解析。2.2 保持 Hook 薄get_conn()只做字段映射AGENTS.md 明确写道PydanticAIHook.get_conn()maps Airflow connection fields to pydantic-ai constructors. That is the hooks entire job. Do not add abstraction layers (builders, factories, registries, Protocols) on top of pydantic-ais own abstractions.源码结构与这条原则完全一致PydanticAIHook的唯一扩展点是_get_provider_kwargs(api_key, base_url, extra)。基类实现只处理最常见的api_keybase_url模式对应连接字段的password与host返回{api_key: ..., base_url: ...}三个云服务子类各自覆写这一个方法即可完成适配没有引入任何注册表、构建器或 ProtocolHook 类conn_type映射的关键字段PydanticAIHookpydanticaipassword→api_keyhost→base_urlextra.model→模型名PydanticAIAzureHookpydanticai_azurehost→azure_endpointextra.api_versionPydanticAIBedrockHookpydanticai_bedrockextra中的region_name、aws_access_key_id/aws_secret_access_key/aws_session_token、profile_name、api_keyBearer Token、base_url、aws_read_timeout/aws_connect_timeout强制转floatPydanticAIVertexHookpydanticai_vertexextra中的project、location、api_key、base_url、service_account_info内联 JSON 对象懒加载google.oauth2.service_account构建凭据这些类的定义见 hooks/pydantic_ai.py连接表单字段隐藏字段、字段重命名、占位提示则由各类的get_ui_field_behaviour()声明并与 provider.yaml 中connection-types段的ui-field-behaviour保持一致。值得注意的细节Bedrock/Vertex 的 UI 会隐藏host/password字段所有配置都存进extra而PydanticAIVertexHook会忽略历史遗留的vertexai字段并打警告——从源码注释看选择 Vertex AI 还是 Generative Language API 现在完全由模型前缀google-cloud:对比google:决定转发旧字段反而会被get_conn()的except TypeError吞掉、静默以错误身份认证。2.3 不提前抽象AGENTS.md 规定不要为单一代码路径添加 Protocol、建造者模式或插件系统等到出现 3 个以上具体用例再引入抽象。对照 src/airflow/providers/common/ai 的实际结构可以验证这一点Hook 只有 5 个文件hooks/pydantic_ai、mcp、langchain、llamaindex及__init__没有registry.py、factory.py之类的基础设施文件每个连接类型一个 Hook 类每个 Operator 一个文件。抽象是等真实需求出现后再长出来的。2.4 Operator 各司其职每个 Operator 只做一件事AGENTS.md 列举的对应关系在 operators/ 目录中一一对应LLMOperator—— prompt → outputllm.pyLLMBranchOperator—— prompt → 分支决策llm_branch.pyLLMSQLOperator—— prompt → 经过校验的 SQLllm_sql.py。provider.yaml 的operators段还登记了agentAgentOperator、llm_file_analysis、llm_schema_compare、document_loader以及 LlamaIndex 的llamaindex_embedding/llamaindex_retrieval与task-decorators段声明的task.agent、task.llm、task.llm_sql等装饰器decorators/构成 Operator/Task 双入口。2.5 一个后端对应一个 ToolsetAGENTS.md 规定Toolset 只封装单一执行后端如DbApiHook、DataFusionEngine若新后端需要相同的工具接口就新建一个 Toolset 类而不是在构造函数里加互斥参数、然后在每个方法里分支判断。用户在 DAG 中通过AgentOperator(toolsets[...])自由组合多个 Toolset。源码印证了这一点toolsets/ 下按后端拆分出sql.pyDbApiHook、datafusion.pyDataFusionEngine、sandbox.py、skills.py、mcp.py、langchain_bridge.py、managed_agent.py、logging.py、hook.py等文件每个文件一个后端。AgentOperator 的构造函数接收llm_conn_id与toolsets: list[AbstractToolset]在durableTrue时还会把每个 Toolset 包进CachingToolset实现步骤级缓存operators/agent.py 的_build_durable_toolsets缓存与不缓存的逻辑同样以包裹而非分支的方式实现。三、新增 LLM Provider 支持三步标准流程AGENTS.md 将新增一个 LLM Provider拆成了明确的条件分支这也是该文档对贡献者最有操作价值的部分。3.1 pydantic-ai 已支持该供应商 → 本包什么都不做用户在 Airflow 连接中直接填写provider:model字符串即可例如azure:gpt-4o、bedrock:anthropic.claude-sonnet-4-20250514Hook 会通过infer_model()自动解析。模型名的取值优先级在get_conn()中是hook 的model_id参数 连接extra里的model字段两者都缺省时抛出ValueError提示Set model_id on the hook or the Model field on the connection。单元测试 test_pydantic_ai.py 对这三条路径逐一验证test_model_id_param_overrides_extramodel_id优先于extra、test_get_conn_with_model_from_extra从extra取anthropic:claude-opus-4-6、test_get_conn_raises_when_no_model无模型名报错、test_get_conn_without_credentials_uses_default_provider无凭据时走环境变量默认路径断言infer_model以bedrock:us.anthropic.claude-v2被调用且不带provider_factory。3.2 需要api_key/base_url之外的凭据 → 覆写_get_provider_kwargs加分支此时应使用 pydantic-ai 自己的 Provider 类如AzureProvider而不是自己拼 SDK 客户端。仓库中现成的三个范例就是按此流程写出来的AzurePydanticAIAzureHook._get_provider_kwargshost映射为azure_endpointextra.api_version透传BedrockPydanticAIBedrockHook._get_provider_kwargs凭据解析顺序为 ①extra.api_keyBearer Token等价于AWS_BEARER_TOKEN_BEDROCK优先于 ②extra中的 IAM 三件套aws_access_key_idaws_secret_access_key可选aws_session_token两者都不填时 ③ 走 AWS 默认凭据链AWS_PROFILE、实例角色等。源码中还有一个容易踩的坑被显式处理了BedrockProvider要求超时参数是float而连接extra是 JSON、整数会原样传入因此aws_read_timeout/aws_connect_timeout被强制float()转换VertexPydanticAIVertexHook._get_provider_kwargs优先级为 ①extra.service_account_info内联 JSON 对象→ 构建Credentials传入GoogleProvider②extra.api_key③ Application Default Credentials。3.3 更新连接表单文档若引入了新字段需要同步更新连接表单文档docs/connections/下按连接类型分文件如 pydantic_ai.rst、pydantic_ai_azure.rst、pydantic_ai_bedrock.rst、pydantic_ai_vertex.rst并保证provider.yaml中connection-types的conn-fields与之一致。反过来如果 pydantic-ai 尚不支持某供应商AGENTS.md 的结论是向上游 pydantic-ai 提贡献而不是在本包内自建 wrapper。这与 2.1 节的委托原则一脉相承。3.4 一个值得注意的实现细节连接测试不做真实调用PydanticAIHook.test_connection() 只调用get_conn()验证模型串合法、Provider 类能用给定凭据实例化不发起真实 LLM API 调用——源码注释解释了原因真实调用昂贵且会因配额、计费、限流等与连通性无关的因素失败。因此 Airflow UI 中的测试连接按钮给出的是解析层面的验证不是端到端验证。四、新增 Toolset 的四条规范AGENTS.md 的Adding a New Toolset一节给出了四条可执行的工程规范每一条都能在现有 Toolset 源码中找到对应实现SQL 类访问遵循四工具模式。凡是提供 SQL 类访问表、schema、查询的 Toolset应复刻SQLToolset的四个工具list_tables、get_schema、query、check_query。在 toolsets/sql.py 中这四个工具各自有独立的 JSON Schema_LIST_TABLES_SCHEMA、_GET_SCHEMA_SCHEMA、_QUERY_SCHEMA、_CHECK_QUERY_SCHEMAcall_tool方法按名字分派其余名字一律拒绝并提示先用 list_tables 和 get_schema 检查数据库。共享逻辑抽成 helper 而非继承。结果截断utils/query_results.py的build_query_result与DEFAULT_MAX_RESULT_BYTES、SQL 校验utils/sql_validation.py的validate_sql、allowed_tables过滤SQLToolset._enforce_allowed_tables都是包内 helperToolset 之间不靠继承树共享逻辑。构造函数保持聚焦。只接受对应当前后端有意义的参数不要加入某些模式下会被静默忽略的参数。SQLToolset的allowed_tables就是正面例子None默认表示允许所有表设置后list_tables只回显授权表、query/check_query会对解析出的语句执行_enforce_allowed_tables越权表名会被拒绝并提示使用list_tables查看可用表。同步 I/O 的工具必须标记sequentialTrue。从源码搜索可以确认这是全 Toolset 的硬性惯例sql.py共享DbApiHook的同步调用、hook.pyhook methods perform synchronous I/O、sandbox.py所有工具共享一个沙箱、datafusion.py 都显式设置了sequentialTrue。另一个体现成本敏感的实现细节query工具没有直接取全量结果再丢弃而是用 toolsets/sql.py 中的_CappedFetch游处理器只取前 N 行——源码注释对比了common.sql的fetch_all_handler整集拉回 worker 再截断指出按游标取数能让传输成本与模型实际看到的内容成比例。五、安全红线两条绝不AGENTS.md 的 Security 一节只有两条但都是不可协商的硬约束。5.1 禁止从连接 extras 做动态导入Never useimportlib.import_module()on user-provided strings from connection fields.理由是 Airflow 连接含 extras对所有持有连接编辑权限的用户可写如果把 extras 里的字符串当作模块路径importlib.import_module()就为低权限用户制造了代码执行向量。这条约束界定了连接 extras 的语义边界它是配置数据不是代码。5.2 SQL 校验默认开启不得关闭LLMSQLOperator对模型生成的 SQL 做 AST 级校验且默认开启。其底层是 utils/sql_validation.py基于 sqlglot 解析sqlglot.parse(sql, dialect..., error_levelErrorLevel.RAISE)策略上fail-closed——默认只允许单条 SELECT 族语句多语句输入被拒绝因为良性 SELECT 后面可以藏危险操作显式拦截数据修改型 CTEWITH d AS (DELETE FROM t RETURNING *) SELECT * FROM d、SELECT INTO、以及用DESCRIBE/EXPLAIN包裹 DDL/DML 的变体sqlglot 无法识别的函数默认全部拒绝可经allowed_functions白名单放行支持 SQLAlchemy dialect 名到 sqlglot dialect 的归一化resolve_sqlglot_dialect未知方言名会被丢弃而不是让解析崩溃。check_query工具与query工具共用同一套_enforce_allowed_tables检查sql.py对应的单元覆盖在 tests/unit/common/ai/utils/test_sql_validation.py 与 tests/unit/common/ai/toolsets/test_sql.py。六、常见陷阱清单PitfallsAGENTS.md 最后列出的三条 Pitfalls 是对前文原则的反向再强调全部有源码可验证不要直接构造裸 Provider SDK 客户端如openai.AsyncAzureOpenAI。应使用 pydantic-ai 的 Provider 类客户端构造重试、超时、认证细节由 pydantic-ai 内部处理。Hook 中唯一外部 SDK 接触面是PydanticAIVertexHook懒加载的google.oauth2.service_account且仅用于把service_account_info转成 pydantic-ai 可接收的credentials对象。不要为供应商新增连接类型。AGENTS.md 原文说单一pydantic_ai连接类型覆盖所有供应商就当前仓库代码而言标准连接类型是pydanticaihooks/pydantic_ai.py 的conn_typeprovider:model字符串区分供应商。pydanticai_azure/pydanticai_bedrock/pydanticai_vertex三个专用类型是例外而非常态——它们只因为认证机制非标准endpointAPI version、IAM 凭据链、服务账号而存在新增pydanticai_openai、pydanticai_anthropic这类按供应商的连接类型是被明确禁止的方向。SDK 导入必须走 compat 层。全 provider 源码统一使用from airflow.providers.common.compat.sdk import ...BaseHook、Connection等禁止直接from airflow.sdk import ...对 src/airflow/providers/common/ai 目录做导入统计可以看到该写法遍布 hooks、operators、decorators 与示例 DAG是兼容性边界约定。七、关键路径索引从哪里继续读AGENTS.md 结尾的 Key Paths 是进入该 provider 源码的官方入口结合 provider.yaml 的组件登记整理成如下索引均为仓库相对路径关注点路径Hookpydantic-ai 桥接hooks/pydantic_ai.py其余 HookMCP / LangChain / LlamaIndexhooks/Operatorsoperators/Task 装饰器decorators/Toolsetstoolsets/持久化执行durable 缓存durable/单元测试tests/unit/common/ai/文档docs/quickstart、toolsets、security、observability、retry_policies 等组件登记与配置项provider.yaml示例 DAGexample_dags/与薄桥设计直接相关的配置项也在provider.yaml的config: common.ai段声明durable_cache_pathAirflow 3.3 时AgentOperator(durableTrue)的每步缓存 ObjectStorage 路径 3.3 存入 AIP-103 任务状态存储后该选项被忽略、otel_export_enabled为 agent 挂上 OpenTelemetry GenAI span默认关、capture_content是否把 prompt/completion 原文写入 span默认关且强调 Airflow 的 secret masking 不作用于 span 属性。create_agent()中的 observability 接线 正是读取这些配置、在 pydantic-ai 2.x 的agent.instrument属性上应用埋点的实现。八、小结Common AI Provider 的整套规范可以概括为一句话把接入 LLM这件事中所有与供应商有关的复杂度推给 pydantic-aiAirflow 侧只保留连接字段映射、Operator 语义封装和 Toolset 后端适配三样东西。对贡献者而言这份 AGENTS.md 给出了清晰的决策树——新供应商先看 pydantic-ai 是否已支持通常什么都不用写非标凭据只覆写_get_provider_kwargs新 Toolset 按单后端 四工具 SQL 模式 helper 共享 sequentialTrue落地安全上守住连接 extras 不可导入、SQL 校验不可关闭两条红线。仓库中的PydanticAIHook四个子类、SQLToolset与AgentOperator的toolsets组合机制就是这套原则的完整实现样例。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考