ARTICLE DETAIL

资讯详情

深耕编程入门与网站建设的一线实战洞察。

Apache Airflow 3.2:通过 `--queues` 将 Trigger 定向分配到指定 Triggerer 主机(Trigger 队列分配机制全解)

Apache Airflow 3.2:通过 `--queues` 将 Trigger 定向分配到指定 Triggerer 主机(Trigger 队列分配机制全解) Apache Airflow 3.2通过--queues将 Trigger 定向分配到指定 Triggerer 主机Trigger 队列分配机制全解【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflowApache Airflow 3.2.0 为 Triggerer 引入了基于任务队列的 Trigger 定向分配能力启动命令新增--queues选项允许每台 triggerer 主机只消费指定任务队列上产生的 Trigger。本篇将围绕这一特性对应 changelog 条目 59239.feature.rst讲清它的适用场景、启用配置、启动命令用法、各类 Trigger 的队列来源差异并结合仓库源码剖析从 CLI 参数到数据库分配查询的完整调用链以及必须遵守的部署注意事项。读完本文你可以按多租户或权限隔离的形态对触发器集群做精细化调度并理解其底层行为边界。1. 特性定位为什么需要按队列分配 Trigger在默认模式下所有 Triggerer 实例对触发器池是无差别共享的任何 triggerer 都可以领取任何未分配的 Trigger。但在以下场景中管理员希望把 Trigger 限定在特定的一组 triggerer 主机上运行引自官方文档 deferring.rst多租户 Airflow 系统按团队划分每个团队运行独立的一组 triggerer异构权限部署不同 triggerer 主机组配置不同的触发器运行环境例如各自持有不同的云权限cloud permissions。Airflow 3.2.0 的解法就是Trigger 队列分配当任务 defer 时将其所属的**任务队列task queue**透传给新注册的 Trigger 实例启动时带--queues参数的 triggerer 只领取队列匹配的 Trigger。2. 启用步骤与启动命令官方文档给出了两步启用流程见 deferring.rst将配置项[triggerer] queues_enabled设为true。这样任务 defer 时会把其分配到的任务队列传给新注册的 Trigger 实例在目标 triggerer 主机的启动命令上追加--queues逗号分隔的任务队列名。该选项确保 triggerer 只处理来自指定任务队列的任务所 defer 出的 Trigger。例如运行两台 triggerer 主机X 和 Y# triggerer X 启动命令 airflow triggerer --queuesalice,bob # triggerer Y 启动命令 airflow triggerer --queuestest_q此时 triggerer X 只会运行来自任务队列alice或bob的任务所 defer 出的 Triggertriggerer Y 只运行来自队列test_q的 Trigger。2.1 配置项triggerer.queues_enabled定义于 config.yml属性值所在段[triggerer]名称queues_enabled类型boolean默认值False引入版本3.2.0说明设为True时deferred 任务注册 Trigger 会带上其来源任务队列triggerer 可基于队列归属选择性运行 Trigger仅在使用支持任务队列分配的 executor 时有意义对应的 CLI 参数定义在 cli_config.pyARG_QUEUES Arg( (--queues,), typestring_list_type, helpOptional comma-separated list of task queues which the triggerer should consume from., )注意其类型为string_list_type即命令行接受逗号分隔字符串并解析为列表。2.2 启动时的强制校验--queues并非随时可用。triggerer_command.py 中的triggerer()入口会做硬性检查if args.queues and not conf.getboolean(triggerer, queues_enabled, fallbackFalse): raise AirflowConfigException( --queues option may only be used when triggerer.queues_enabled is True. ) ... queues set(args.queues) if args.queues else None也就是说若未开启queues_enabled就传了--queues进程直接以AirflowConfigException失败退出——配置开关与 CLI 选项必须成对出现。校验通过后队列集合queues或None表示不限队列会被传入执行链。3. 源码剖析从 CLI 到 Trigger 分配查询3.1 参数传递链triggerer()校验后将queues传入triggerer_run()triggerer_command.py后者构造TriggererJobRunnertriggerer_job_runner TriggererJobRunner( jobJob(heartratetriggerer_heartrate, team_nameteam_name), capacitycapacity, queuesqueues, )TriggererJobRunner.__init__triggerer_job_runner.py保存self.queues queues随后在_execute()中把队列集合连同capacity、team_name一起交给TriggerRunnerSupervisor.start(...)启动的看门狗/消息桥接层。3.2 核心每一轮循环中的分配与领取TriggerRunnerSupervisor的主循环run_once()triggerer_job_runner.py每个 tick 依次执行加载 Trigger、服务子进程、处理事件、处理失败、清理、心跳、打点。其中load_triggers()是队列过滤生效的关键def load_triggers(self) - None: Assign triggers to this triggerer and update the runner with the IDs it should run. ... Trigger.assign_unassigned( self.job.id, self.capacity, self.health_check_threshold, queuesself.queues, team_nameself.team_name, ) ids Trigger.ids_for_triggerer(self.job.id, queuesself.queues, team_nameself.team_name) self.update_triggers(set(ids))两个方法都在 trigger.py 模型层实现1Trigger.ids_for_triggerertrigger.py决定本机应该运行哪些 Trigger其中的注释直接点明了语义query select(cls.id).where(cls.triggerer_id triggerer_id) # By default, there is no trigger queue assignment. Only filter by queue when explicitly set in the triggerer CLI. # Filter by queues if the triggerer explicitly was called with --queues, otherwise, filter out # Triggers which have an explicit queue value since there may be other triggerer hosts explicitly assigned to that queue. if queues: query query.filter(cls.queue.in_(queues)) else: query query.filter(cls.queue.is_(None))2Trigger.assign_unassignedtrigger.py负责认领新 Trigger先按本机已持有数量扣减 capacity再选出没有存活 triggerer 领取的 Triggerget_sorted_triggers按 Callback → 任务实例 → Asset 的优先级排序最后批量把triggerer_id更新为本机。queues参数同样参与候选过滤。从源码结构看这组过滤逻辑揭示了一条容易被忽视的部署约束与官方文档的 Caveats 完全一致不带--queues的 triggerer 只消费queue为NULL的 Trigger带--queues的 triggerer 只消费queue在其列表内的 Trigger。二者天然互斥、互不重叠。3.3 Trigger 侧的队列继承机制任务 defer 时队列是如何写进 Trigger 的BaseTrigger定义了类属性base.py# Whether a deferred tasks queue should be inherited by this trigger when # triggerer.queues_enabled is set. Subclasses that assign their own queue # directly (e.g. BaseEventTrigger, CallbackTrigger) must set this to False, # or _defer_task will overwrite that queue with the deferring tasks queue. trigger_queue_inherited_from_task: bool True # Trigger queue assignment. None means no explicit assignment; BaseEventTrigger and # CallbackTrigger set this in their __init__ queue: str | None None可以这样理解普通任务创建型Trigger 默认trigger_queue_inherited_from_task True在queues_enabledTrue时由 defer 流程把任务队列继承为queue字段而事件驱动 Trigger 与回调 Trigger 在各自__init__中显式设置queue并关闭继承避免被任务队列覆盖。4. 各 Trigger 类型的队列来源与分配规则官方文档给出了一张权威对照表deferring.rst在queues_enabledTrue时各类 Trigger 的 triggerer 分配来源如下Trigger 类型Triggerer 分配来源queues_enabled为 True 时任务创建的 Trigger 实例其--queues中包含该任务队列的任意 triggerer事件驱动 TriggerEvent-Driven Triggers--queues包含该 Triggerqueue的任意 triggerer若未设置queue则任何不带--queues的 triggerer 均可来自异步 Callback 的 Trigger--queues包含该 callbackqueue的任意 triggerer若未设置queue则任何不带--queues的 triggerer 均可补充说明同样来自上述文档该特性对任务创建型 Trigger仅兼容利用任务queue概念的 executor如 CeleryExecutor事件驱动 Trigger 与任务队列无关但BaseEventTrigger子类可以在super().__init__()中显式传queue指定归属队列异步 Callback 的 Trigger 可通过给AsyncCallback传queue指定归属队列事件驱动 Trigger 的queue在 DAG 处理阶段注册AssetWatcher时一次性确定callback Trigger 的queue在回调被投递到 Triggerer 时确定。两者都与 Multi-Team 模式的team_name作用域相互独立、可以组合使用。4.1 多 triggerer HA 混部示例文档给出了一个所有类型 Trigger 都能被运行的 HA 混部示例假设所有任务的 queue 值只取team_A/team_B见 deferring.rst# triggerer A 启动命令仅消费由 queue team_A 中任务注册的 triggers airflow triggerer --queuesteam_A # triggerer B 启动命令仅消费由 queue team_B 中任务注册的 triggers airflow triggerer --queuesteam_B # triggerer C 启动命令仅消费事件型与 callback 型 triggers airflow triggerer其中不带--queues的 C 承担了queue IS NULL的兜底角色。5. 必须遵守的部署注意事项Caveats以下约束直接决定该功能是否可用、Trigger 是否会卡死无人消费executor 兼容性仅对支持任务队列分配的 executor 生效对任务创建型 Trigger必须保留兜底 triggerer如果你同时运行有队列与无队列的 Trigger任意类型必须至少运行一台不带--queues的 triggerer 主机否则后者永远不会被运行队列名要能对上若把某台 triggerer 的--queues设成没有任何 Trigger 归属的队列值该 triggerer 将永远不运行任何 Trigger。反过来所有不带--queues的 triggerer 只会消费未显式设置queue的 callback 型与事件驱动型 Trigger与 Multi-Team 的关系--queues是独立的队列机制若使用 Multi-Team 模式--team-name提供按团队的原生作用域分配两者可叠加使用--team-name要求core.multi_team开启同样在 triggerer_command.py 中做存在性校验。6. 测试验证与相关配置参考该特性的行为由单测系统覆盖CLI 校验与队列透传test_triggerer_command.py使用conf_vars({(triggerer, queues_enabled): True})验证--queues路径分配/查询逻辑test_trigger.py 中专门有当triggerer.queues_enabled设为True时验证 trigger 分配处理的测试类约 651 行起并对queues_enabled为 True/False 两种情形分别断言 Trigger 注册时queue字段取值如 554、752、838 行附近的参数化用例Execution API 侧对 deferred 任务queue字段的返回也有针对queues_enabled双态的参数化测试test_task_instances.py其中可见启用后队列值会落到默认为default的queue字段。此外同属[triggerer]段、与 HA 负载均衡相关的max_trigger_to_select_per_loop默认 503.2.0 引入见 config.yml与本特性属于同一版本波次它限制每轮循环最多选中的 Trigger 数避免多个 triggerer 争抢时饿死彼此可与队列分配组合使用。7. 小结要点说明特性airflow triggerer --queuesq1,q2使 triggerer 只消费指定任务队列的 Trigger前置条件[triggerer] queues_enabled True默认False3.2.0 引入否则 CLI 启动即报AirflowConfigException过滤规则带--queues的 triggerer 只匹配queue IN (...)不带的只匹配queue IS NULL队列来源任务型 Trigger 继承任务队列事件驱动/Callback 型 Trigger 可显式指定queue部署约束混部有/无队列Trigger 时必须保留至少一台不带--queues的兜底 triggerer关键源码triggerer_command.py、triggerer_job_runner.py、trigger.py 中的assign_unassigned/ids_for_triggerer、base.py 中的queue/trigger_queue_inherited_from_task至此你既掌握了该特性面向多租户与权限隔离场景的完整操作路径也理解了配置开关 → CLI 参数 → JobRunner → 模型层 SQL 过滤这条源码链路以及混部部署中避免 Trigger 无人消费的兜底原则。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表