
Apache Airflow 3.3 Worker-Dispatched 异步连接测试从 API 服务器迁移到 Worker 执行的完整指南【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow导读Apache Airflow 在 3.3.0 中引入了一项关键特性worker-dispatched异步连接测试。传统上POST /connections/test同步端点会在 API 服务器进程内直接执行连接测试这意味着测试所用的凭据与网络路径都暴露在 API 服务器上。而新的异步工作流允许连接测试通过执行器executor被派发到worker上运行——与真实任务运行的位置完全一致——并通过POST /connections/enqueue-test提交、以 token 轮询结果。本文将基于本仓库源码完整讲解该特性的启用方式、REST API 用法、[connection_test]配置项、底层状态机与安全模型帮助你理解并落地这一能力。一、为什么需要把连接测试放到 Worker 上执行Airflow 的Connection对象用于存储连接外部服务所需的凭据与其他信息参见 connections 概念文档。在引入异步测试之前连接测试只有一个路径同步测试调用POST /connections/test测试在 API 服务器webserver 进程内同步执行。这种模式存在两个现实痛点官方文档明确指出了它们凭据在 API 服务器上被使用连接的用户名/密码等敏感字段会被 API 服务器进程读取并用于建立真实连接而你或许并不希望凭据在 API 服务器上被演练网络可达性差异很多连接例如生产数据库只允许从 worker 网络段访问API 服务器根本连不通它们。异步worker-dispatched连接测试正是针对这两个问题设计的测试在 worker 上执行凭据只在 worker 上被使用——也就是你的任务真正运行的地方而非 API 服务器。注意这一特性与同步测试共用同一个总开关——[core] test_connection配置项。也就是说无论哪种测试方式都必须先启用连接测试功能详见下文。二、特性总览两个端点、一个 token、三种配置该特性源自 newsfragment 62343.feature.rst的核心交互模型非常简单步骤操作说明1POST /connections/enqueue-test提交连接测试请求返回一个轮询 tokenHTTP 202 Accepted2GET /connections/enqueue-test携带Airflow-Connection-Test-Token请求头轮询测试状态与结果3可选的收尾若测试成功且请求了commit_on_success连接会被写入元数据库的connection表同时原有的同步端点POST /connections/test保持不变两种模式可以并存由调用方按需选择。2.1 启用前置条件[core] test_connection无论同步还是异步连接测试默认都是关闭的这是出于安全考虑——官方强烈建议仅在确保只有高度可信的 UI/API 用户拥有 edit connection 权限时才启用。开关由配置项[core] test_connection控制可用环境变量AIRFLOW__CORE__TEST_CONNECTION覆盖取值有三种Disabled禁用测试功能同时禁用 UI 中的 Test Connection 按钮默认值Enabled启用测试功能并激活 UI 按钮Hidden禁用测试功能并隐藏 UI 按钮。在源码层面这一开关由路由层的_ensure_test_connection_enabled()强制执行若配置值不是enabled直接返回 HTTP 403Testing connections is disabled in Airflow configuration...见 connections.py。POST /connections/test与POST /connections/enqueue-test两个端点都会先调用这个守卫函数。2.2 配置段[connection_test]异步测试的行为由新增的[connection_test]配置段version_added: 3.3.0调控共三个参数定义于 config.yml配置项类型默认值作用connection_test.timeoutinteger60一次 worker-dispatched 测试允许运行的最大秒数超过即视为超时调度器 reaper 会使用该值加上宽限期grace period将陈旧测试标记为 failedconnection_test.max_concurrencyinteger4同时处于活跃状态QUEUED RUNNING的测试数量上限超出的请求停留在 PENDING 状态等待槽位connection_test.reaper_intervalfloat30.0调度器每隔多少秒检查一次陈旧测试QUEUED/RUNNING 超过 timeout 宽限期并将其标记为失败一个值得注意的实现细节max_concurrency上限是按调度器per-scheduler而非全局强制的。在有 N 个 HA 调度器的部署中最坏情况下每次 tick 的派发量为N * max_concurrency。由于连接测试是用户发起的、发生频率低超出部分会通过 reaper 自行纠正——这是 config.yml 中明确说明的设计取舍。三、REST API 实战提交与轮询3.1 提交测试POST /connections/enqueue-test提交请求体为ConnectionTestRequestBody包含连接定义的全部字段并可附加三个控制字段commit_on_success测试成功时是否将连接写入或更新真实的connection表executor可选的执行器名称必须属于已配置的执行器列表否则返回 422见_ensure_executor_is_configured()校验逻辑queue可选的队列名。curl -X POST https://airflow-host/api/v1/connections/enqueue-test \ -H Authorization: Bearer token \ -H Content-Type: application/json \ -d { connection_id: my_prod_db, conn_type: postgres, login: my-login, password: my-password, host: db.internal.example.com, port: 5432, schema: analytics, commit_on_success: false, executor: kubernetes }成功时返回HTTP 202 Accepted响应体包含token、connection_id与初始state{ token: x8K...3Q, connection_id: my_prod_db, state: pending }在提交阶段路由还会做两件安全校验见 enqueue_connection_test 实现团队一致性若该 connection_id 在元数据库中已存在且带有团队归属请求中的team_name必须与之匹配否则返回 403授权检查调用 auth manager 的is_authorized_connection校验当前用户对目标连接是否具备测试权限按是否commit_on_success区分要求PUT或POST方法权限。此外若同一 connection_id 已存在活跃的测试请求插入将触发唯一约束冲突uq_connection_test_request_active_connAPI 返回 409 ConflictAn active connection test already exists...。3.2 轮询结果GET /connections/enqueue-test用上一步拿到的 token 作为请求头轮询curl https://airflow-host/api/v1/connections/enqueue-test \ -H Authorization: Bearer token \ -H Airflow-Connection-Test-Token: x8K...3Q响应体包含当前state、result_message与创建时间等字段。state取值遵循完整的生命周期状态机见下节。安全要点结果只能通过该 token 读取且仅对获得该连接授权的用户可见在 multi-team多团队部署中其他团队发起的连接测试对你完全不可见——如果 token 不存在或用户无权访问接口统一返回 404No connection test found for token...避免泄露连接 ID 的存在性见 get_connection_test 实现。四、底层实现剖析状态机、持久化与执行链路4.1connection_test_request表与状态机异步测试请求被持久化到新表connection_test_request对应 ORM 模型ConnectionTestRequest见 connection_test.py迁移脚本见 0118_3_3_0_add_connection_test_table.py。该表完整存储连接字段conn_type、host、login、password、schema、port、extra使 worker 从这张表读取测试数据而不触碰真实的connection表只有当测试成功且commit_on_successTrue时commit_to_connection_table()才会对真实连接表执行 upsert代码见 L205-L230敏感字段password、extra经 Fernet 加密存储支持密钥轮换通过tokensecrets.token_urlsafe(32)生成唯一标识请求通过active_connection_id与唯一约束保证同一连接同时只有一个活跃测试。状态机定义在ConnectionTestState枚举L52-L62pending → queued → running → success / failedACTIVE_STATES {pending, queued, running}活跃期间active_connection_id被同步设置TERMINAL_STATES {success, failed}4.2 完整执行链路从提交到结果落库一次异步测试的调用链如下API 服务器POST /connections/enqueue-test创建ConnectionTestRequest初始 statepending返回 token调度器以固定节奏扫描 PENDING 测试结合max_concurrency槽位限制将其派发QUEUED到指定 executor 的队列connection-test workload 的定义见 executors/workloads/connection_test.pyWorker通过 Execution API 的GET /{connection_test_id}/connection一次性取回连接数据with_for_updateTrue行锁 原子地将状态从 pending/queued 翻转为 running且凭据只能被取一次重复获取返回 409见 connection_tests.py随后在 worker 上调用 hook 的test_connection方法执行真实测试Worker通过PATCH /{connection_test_id}回报结果state result_message若为 success 且请求了commit_on_success在此处调用commit_to_connection_table()写回连接表L139-L140reaper调度器每隔reaper_interval秒检查一次将超过timeout含宽限期仍处于 queued/running 的测试标记为 failed调用方用 token 轮询GET /connections/enqueue-test读取最终状态与结果消息。4.3 双轨安全模型轮询 token 与 worker JWT这是本特性最值得关注的安全设计官方文档connection.rst 的异步测试章节与源码共同佐证面向调用方的轮询 token只用于从 API 读取结果绑定到特定连接与团队无权限者 404面向 worker 的执行授权是独立的一层调度器为单次请求签发短期 JWT其subject 就是 connection-test request idscope 为workload。Execution API 的端点强制ct:self检查路由级依赖Security(require_auth, scopes[ct:self, token:workload])见 connection_tests.py L37-L42——要求 token 的 subject 与请求路径中的 connection-test id完全匹配。这意味着worker 拿到的 token 只能取回这一个测试请求的连接数据并回报其结果无法触达其他连接测试、任务实例task instances或任何连接。即使 worker 被攻破攻击面也被严格限制在单次请求范围内。五、适用场景与注意事项推荐使用异步测试的场景连接只能从 worker 网络段访问API 服务器无法直连不希望凭据在 API 服务器上被使用/演练希望凭据只在 worker 上被接触。需要注意的几点官方文档明确提示结果取决于 worker 环境由于测试在 worker 上执行其结果反映的是该 worker 上安装的库、providers 与网络访问能力可能与 API 服务器不同。若 worker 与 API 服务器或不同 worker 之间安装的库/provider 版本不一致测试结果可能产生差异不可用于外部 secrets backend 中的连接与同步测试一致存放在外部 secrets backend如 Vault、AWS SSM中的连接通过 UI 或 REST API 无法使用该测试功能依赖 hook 实现Airflow 通过调用连接类型对应 hook 类的test_connection方法来执行测试若该连接类型没有关联 hook或 hook 未实现test_connection将返回错误信息凭据加密存储异步请求中的密码等敏感字段在元数据库中以 Fernet 加密存储is_encrypted/is_extra_encrypted标记与连接表的安全策略一致。六、验证与测试参考仓库单元测试覆盖了本特性的完整行为见 test_connections.py包括token 轮询、401/403/404/409 等异常路径、团队授权边界、commit_on_success落库等。若你想在本地验证异步测试核心操作路径可归纳为三步airflow config确认[core] test_connection Enabled→POST /connections/enqueue-test获取 token → 携带Airflow-Connection-Test-Token头轮询GET /connections/enqueue-test直至 state 到达success或failed终态。OpenAPI 规范v2-rest-api-generated.yaml与 UI 类型定义services.gen.ts中均可查看到这两个端点的完整请求/响应结构可作为集成开发的权威参考。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考