
Apache Airflow 手动终结 Dag Run 时触发任务实例监听器Bugfix 69874 行为深度解析【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow当运维人员通过 REST API 或 Web UI 手动将 Dag Run 状态置为终态success / failed时那些正在运行的任务实例会被强制终结——但此前挂在它们身上的任务实例级监听器on_task_instance_success/on_task_instance_failed却不会被触发导致监控、通知与审计逻辑在人工干预场景下静默漏报。本篇文章以 Apache Airflow 当前仓库中的 69874.bugfix.rst 为核心结合源码与单元测试完整解析该缺陷的修复语义、底层调用链、teardown 任务的特殊处理以及如何编写可感知手动状态变更的监听器插件。1. Bugfix 声明原文与问题背景本次修复对应的变更声明airflow-core/newsfragments/69874.bugfix.rst内容如下Listeners registered viaon_task_instance_success/on_task_instance_failedare now called for non-teardown task instances that were running when a Dag Run state is manually set to a terminal state (e.g. via the API or UI).翻译并拆解其语义可以提炼出四个关键信息触发场景Dag Run 状态被手动设置为终态例如通过 API 或 UI而不是由调度器/任务执行自然流转到终态作用对象当时处于运行状态的任务实例running 类状态排除对象teardown 任务被有意排除在外触发结果on_task_instance_success/on_task_instance_failed监听器现在会被正确调用。修复前的典型痛点一个 DAG 中某任务正在长耗时运行人工从 UI 点击标记成功/标记失败将整个 Dag Run 强制终结。虽然数据库里该任务实例的状态被改写为终态但依赖该状态变更的监听器如发 Slack 通知、上报指标、写审计日志毫无感知形成状态已变、事件未发的割裂。2. 修复后的行为模型修复后手动终结 Dag Run 的完整事件模型如下表手动设置的 Dag Run 终态运行中任务实例的处理触发的任务实例级监听器触发的 DAG Run 级监听器success置为SUCCESS强制终结on_task_instance_successon_dag_run_successfailed置为FAILED强制终结on_task_instance_failed带 error 信息on_dag_run_failedqueued不入终态逻辑不触发不触发见下文说明需要补充的两个细节非运行中任务不触发处于 pending未完成状态的任务实例会被置为SKIPPED但不会触发任务实例级监听器——监听器事件只针对正在运行却被强制终结的实例这与on_task_instance_skipped的既有语义保持一致该 hook 只覆盖任务自行跳过详见下文 hookspec 说明queued 状态不通知手动置为queued时只调用set_dag_run_state_to_queued代码注释明确说明Not notifying on queued - only notifying on RUNNING, which happens in the scheduler即 queued 是过渡态而非终态不产生监听器事件。这里的运行状态是广义的活跃状态集合包含四种RUNNING、DEFERRED延迟、UP_FOR_RESCHEDULE等待重调度、AWAITING_INPUT等待输入。换句话说一个处于 deferrable 模式挂起、或等待 sensor 输入的任务实例也会被视为正在运行并进入强制终结 监听器触发流程。3. 底层调用链从 REST 路由到监听器 Hook要理解修复的实现需要沿着一次PATCH /dags/{dag_id}/dagRuns/{dag_run_id}请求走完整个调用链。UI 上的标记成功/失败操作最终同样走这个 REST 端点因此 API 与 UI 两条入口共享同一套逻辑。3.1 路由层patch_dag_run路由定义在 dag_run.py 附近。值得注意的一个细节是 L238-L244 的注释与顺序处理# Apply note before state so listeners fired inside patch_dag_run_state() see the updated note.即先应用 note 再应用 state确保在patch_dag_run_state()内部触发的监听器能够看到本次请求同时更新的备注字段。3.2 服务层patch_dag_run_state核心入口是 services/public/dag_run.py 中的patch_dag_run_state其执行序列为调用set_dag_run_state_to_success/set_dag_run_state_to_failed来自airflow.api.common.mark_tasks得到元组(all_updated_tis, killed_tis)对killed_tis调用_emit_state_listener_hooks(killed_tis, TaskInstanceState.SUCCESS/FAILED)——这是本次修复的关键新增逻辑调用 DAG Run 级监听器on_dag_run_success/on_dag_run_failed消息为Dag Runs state was manually set to success.之类并保证dag_run.dag已挂载方便监听器访问 DAG 信息所有监听器调用均包裹在 try/except 中监听器抛出的异常只记录日志log.exception(error calling listener)不会影响状态变更本身。3.3 状态变更层_set_dag_run_terminal_stateset_dag_run_state_to_success/set_dag_run_state_to_failed都委托给 mark_tasks.py 中的_set_dag_run_terminal_state。该函数的行为筛选出处于四种活跃状态RUNNING、DEFERRED、UP_FOR_RESCHEDULE、AWAITING_INPUT的任务实例不杀 teardown 任务# Do not kill teardown tasks即 teardown 任务被排除在强制终结名单之外不跳过 teardown 任务# Do not skip teardown tasks即 pending 的 teardown 也不被置为 SKIPPED只有当该 Dag Run 中不存在任何 pending teardown 时才把 Dag Run 本身置为终态返回(all_updated_tis, killed_tis)其中killed_tis只包含非 teardown、且处于活跃运行状态、被强制终结的任务实例。函数 docstring 对killed_tis的语义解释得非常清楚mark_tasks.pykilled_tiscontains only the non-teardown TIs that were in an active running state and were forcefully terminated (teardown TIs are intentionally left running so they can finish their own cleanup, and must not receive a terminal listener event here).这一设计决定了后续监听器事件的边界只有killed_tis会收到任务实例级监听器事件。3.4 事件发射层_emit_state_listener_hooks真正触发任务实例级监听器的是 services/public/task_instances.py 中的_emit_state_listener_hooks其逻辑按新状态分发SUCCESS→on_task_instance_success(previous_stateNone, task_instanceti)FAILED→on_task_instance_failed(previous_stateNone, task_instanceti, errorTaskInstances state was manually set tofailed.)SKIPPED→on_task_instance_skipped(previous_stateNone, task_instanceti)每个ti的调用同样被 try/except 包裹单个监听器出错只记日志不影响后续实例的遍历。4. 为什么 teardown 任务被排除这是本次修复语义中最容易误解的一点。teardown 任务在 Airflow 的 setup/teardown 机制中承担资源清理职责例如关闭集群、释放锁、删除临时资源它们必须在主任务被强制终结后继续运行来完成清理而不是被一并杀死。因此手动终结 Dag Run 时teardown 任务保持原状继续执行由 worker/triggerer 正常驱动其完成teardown 任务不会在此过程中收到on_task_instance_success/on_task_instance_failed事件——它们自身状态的最终变化由执行端自然上报走正常执行流程的事件路径而非手动状态变更路径只有当 Dag Run 中还有 pending teardown 时Dag Run 自身不会立即被置为终态mark_tasks.py以保证 teardown 有机会被调度执行。这一设计在 mark_tasks.py 和 L300-L301 中有直接的代码佐证Do not kill teardown tasks、Do not skip teardown tasks。仓库中的示例 DAG example_setup_teardown.py 和 example_setup_teardown_taskflow.py 展示了如何通过as_teardown(setups...)或teardown装饰器标记这类任务。5. 监听器 Hook 签名与参数语义任务实例级监听器的规范定义位于 shared/listeners/src/airflow_shared/listeners/spec/taskinstance.py相关签名如下hookspec def on_task_instance_running(previous_state, task_instance): ... hookspec def on_task_instance_success(previous_state, task_instance): ... hookspec def on_task_instance_failed(previous_state, task_instance, error): ...需要注意的previous_state参数在手动状态变更场景下由于 API 服务器上的TaskInstance快照不携带变更前的状态信息因此previous_state固定为None。这与正常执行流程中由执行端传入真实 previous state的行为不同监听器代码应据此区分事件来源。示例插件 event_listener.py 中正是用isinstance(task_instance, TaskInstance)判断事件是否来自 API 路径hookimpl def on_task_instance_success(previous_state, task_instance): print(Task instance in success state) print( Previous state of the Task instance:, previous_state) if isinstance(task_instance, TaskInstance): print(Task instances state was changed through the API.) print(fTask operator:{task_instance.operator}) return context task_instance.get_template_context() operator context[task] print(fTask operator:{operator})上述hookimpl与hookspec基于 pluggy 插件机制监听器管理器在 airflow-core/src/airflow/listeners/listener.py 中通过get_listener_manager()统一装配并通过integrate_listener_plugins加载用户插件。DAG Run 级 hookon_dag_run_success/on_dag_run_failed的规范定义见 airflow-core/src/airflow/listeners/spec/dagrun.py。6. 编写一个感知手动终结的监听器插件结合上述语义可以编写一个同时关注 Dag Run 级与任务实例级事件的监听器插件以示例插件 event_listener.py 为蓝本将文件放入$AIRFLOW_HOME/plugins目录即可被加载from airflow.listeners import hookimpl from airflow.models.taskinstance import TaskInstance from airflow.utils.state import TaskInstanceState hookimpl def on_task_instance_success(previous_state, task_instance): # 手动终结场景previous_state 为 None且 task_instance 是 API 服务器上的 TaskInstance if previous_state is None and isinstance(task_instance, TaskInstance): print(f[manual-terminate] task {task_instance.task_id} f(dag{task_instance.dag_id}, run{task_instance.run_id}) marked SUCCESS via API/UI) # 正常执行流程previous_state 为真实前置状态 else: print(f[normal] task {task_instance.task_id} succeeded, previous_state{previous_state}) hookimpl def on_task_instance_failed(previous_state, task_instance, error): print(f[manual-terminate] task {task_instance.task_id} marked FAILED: {error}) hookimpl def on_dag_run_failed(dag_run, msg): print(fDAG run {dag_run.dag_id}/{dag_run.run_id} failed: {msg})判定技巧小结previous_state is Noneisinstance(task_instance, TaskInstance)组合基本可以认定事件来自手动状态变更API 路径正常执行路径下task_instance多为RuntimeTaskInstance见 shared/listeners/src/airflow_shared/listeners/spec/taskinstance.py 的类型注解可通过get_template_context()拿到完整渲染上下文on_task_instance_failed的error参数在手动场景下固定为TaskInstances state was manually set tofailed.可用于与真实执行异常区分。7. 单元测试如何验证该修复本次修复行为在仓库单元测试中有完整覆盖测试位于 test_dag_run.py测试用例验证点test_patch_dag_run_notifies_ti_listeners_for_running_tasksL1759运行中的任务实例收到任务实例级监听器事件随后收到 DAG Run 级事件事件顺序与状态符合预期listener.state[0]为 TI 状态、listener.state[1]为 DagRun 状态test_patch_dag_run_does_not_notify_ti_listeners_for_non_running_tasksL1795非运行中queued任务实例不触发任务实例级监听器只有 DAG Run 级事件test_patch_dag_run_does_not_notify_ti_listeners_for_running_teardown_tasksL1829运行中的 teardown 任务被有意跳过不产生任务实例级监听器事件test_patch_dag_run_listener_sees_note_when_note_and_state_both_patchedL1739同一请求同时修改 note 与 state 时监听器能看到更新后的 note验证 3.1 节的顺序处理此外任务实例单点 PATCHPATCH /dags/{dag_id}/dagRuns/{dag_run_id}/taskInstances/{task_id}路径同样有对应测试test_patch_task_instance_notifies_listenerstest_task_instances.py它验证的是同一_emit_state_listener_hooks函数在单实例场景下的行为。这些测试共同构成了本次修复的回归保障。8. 实践建议与注意事项适用场景基于监听器做实时告警/通知Slack、钉钉、邮件确保人工标记成功/标记失败也会发出通知成本核算与资源清理审计统计哪些运行中的任务被人工强制终结避免长时间挂起任务产生费用盲区元数据采集/数据同步监听器将任务状态变化写入外部系统时人工干预同样需要同步。注意事项监听器只对非 teardown、处于活跃运行状态RUNNING / DEFERRED / UP_FOR_RESCHEDULE / AWAITING_INPUT的任务实例生效已终态、queued、pending 的任务不会收到任务实例级事件teardown 任务的清理逻辑不会走这条事件路径如需感知其最终结果请监听其正常执行流程的事件手动场景下previous_state恒为None不要基于它推断前置状态错误消息字符串为固定文案也不要把error当真实异常对象处理虽然签名允许None | str | BaseException监听器内部抛出异常会被吞掉并记录日志不影响状态变更的提交因此请勿把关键业务逻辑放在监听器异常后的补救里本修复针对的是手动设置 Dag Run 终态这一入口调度器自然完成 Dag Run 的执行路径本就通过执行端触发监听器行为不受影响。通过理解 Bugfix 69874 的完整语义你可以在编写 Airflow 监听器插件时准确区分自然流转与人工干预两种事件来源构建更可靠、无遗漏的监控与自动化体系。【免费下载链接】airflowApache Airflow - A platform to programmatically author, schedule, and monitor workflows项目地址: https://gitcode.com/GitHub_Trending/ai/airflow创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考