Apache Airflow 手动终结 Dag Run 时触发任务实例监听器:Bugfix 69874 行为深度解析
2026/9/10 1:22:13 网站建设 项目流程

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).

翻译并拆解其语义,可以提炼出四个关键信息:

  1. 触发场景:Dag Run 状态被手动设置为终态(例如通过 API 或 UI),而不是由调度器/任务执行自然流转到终态;
  2. 作用对象:当时处于运行状态的任务实例(running 类状态);
  3. 排除对象teardown 任务被有意排除在外;
  4. 触发结果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_success
failed置为FAILED(强制终结)on_task_instance_failed(带 error 信息)on_dag_run_failed
queued不入终态逻辑不触发不触发(见下文说明)

需要补充的两个细节:

  • 非运行中任务不触发:处于 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 是过渡态而非终态,不产生监听器事件。

这里的"运行状态"是广义的活跃状态集合,包含四种:RUNNINGDEFERRED(延迟)、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,其执行序列为:

  1. 调用set_dag_run_state_to_success/set_dag_run_state_to_failed(来自airflow.api.common.mark_tasks),得到元组(all_updated_tis, killed_tis)
  2. killed_tis调用_emit_state_listener_hooks(killed_tis, TaskInstanceState.SUCCESS/FAILED)——这是本次修复的关键新增逻辑;
  3. 调用 DAG Run 级监听器on_dag_run_success/on_dag_run_failed,消息为"Dag Run's state was manually set to 'success'."之类,并保证dag_run.dag已挂载,方便监听器访问 DAG 信息;
  4. 所有监听器调用均包裹在 try/except 中,监听器抛出的异常只记录日志(log.exception("error calling listener")),不会影响状态变更本身。

3.3 状态变更层:_set_dag_run_terminal_state

set_dag_run_state_to_success/set_dag_run_state_to_failed都委托给 mark_tasks.py 中的_set_dag_run_terminal_state。该函数的行为:

  • 筛选出处于四种活跃状态(RUNNINGDEFERREDUP_FOR_RESCHEDULEAWAITING_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.py):

killed_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,其逻辑按新状态分发:

  • SUCCESSon_task_instance_success(previous_state=None, task_instance=ti)
  • FAILEDon_task_instance_failed(previous_state=None, task_instance=ti, error="TaskInstance's state was manually set tofailed.")
  • SKIPPEDon_task_instance_skipped(previous_state=None, task_instance=ti)

每个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 instance's state was changed through the API.") print(f"Task operator:{task_instance.operator}") return context = task_instance.get_template_context() operator = context["task"] print(f"Task operator:{operator}")

上述@hookimpl@hookspec基于 pluggy 插件机制,监听器管理器在 airflow-core/src/airflow/listeners/listener.py 中通过get_listener_manager()统一装配,并通过integrate_listener_plugins加载用户插件。DAG Run 级 hook(on_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(f"DAG run {dag_run.dag_id}/{dag_run.run_id} failed: {msg}")

判定技巧小结:

  • previous_state is None+isinstance(task_instance, TaskInstance)组合,基本可以认定事件来自手动状态变更(API 路径);
  • 正常执行路径下task_instance多为RuntimeTaskInstance(见 shared/listeners/src/airflow_shared/listeners/spec/taskinstance.py 的类型注解),可通过get_template_context()拿到完整渲染上下文;
  • on_task_instance_failederror参数在手动场景下固定为"TaskInstance's state was manually set tofailed.",可用于与真实执行异常区分。

7. 单元测试如何验证该修复

本次修复行为在仓库单元测试中有完整覆盖,测试位于 test_dag_run.py:

测试用例验证点
test_patch_dag_run_notifies_ti_listeners_for_running_tasks(L1759)运行中的任务实例收到任务实例级监听器事件,随后收到 DAG Run 级事件,事件顺序与状态符合预期(listener.state[0]为 TI 状态、listener.state[1]为 DagRun 状态)
test_patch_dag_run_does_not_notify_ti_listeners_for_non_running_tasks(L1795)非运行中(queued)任务实例触发任务实例级监听器,只有 DAG Run 级事件
test_patch_dag_run_does_not_notify_ti_listeners_for_running_teardown_tasks(L1829)运行中的 teardown 任务被有意跳过,不产生任务实例级监听器事件
test_patch_dag_run_listener_sees_note_when_note_and_state_both_patched(L1739)同一请求同时修改 note 与 state 时,监听器能看到更新后的 note(验证 3.1 节的顺序处理)

此外,任务实例单点 PATCH(PATCH /dags/{dag_id}/dagRuns/{dag_run_id}/taskInstances/{task_id})路径同样有对应测试test_patch_task_instance_notifies_listeners(test_task_instances.py),它验证的是同一_emit_state_listener_hooks函数在单实例场景下的行为。这些测试共同构成了本次修复的回归保障。

8. 实践建议与注意事项

适用场景

  • 基于监听器做实时告警/通知(Slack、钉钉、邮件),确保人工"标记成功/标记失败"也会发出通知;
  • 成本核算与资源清理审计:统计哪些运行中的任务被人工强制终结,避免长时间挂起任务产生费用盲区;
  • 元数据采集/数据同步:监听器将任务状态变化写入外部系统时,人工干预同样需要同步。

注意事项

  1. 监听器只对非 teardown、处于活跃运行状态(RUNNING / DEFERRED / UP_FOR_RESCHEDULE / AWAITING_INPUT)的任务实例生效;已终态、queued、pending 的任务不会收到任务实例级事件;
  2. teardown 任务的清理逻辑不会走这条事件路径,如需感知其最终结果,请监听其正常执行流程的事件;
  3. 手动场景下previous_state恒为None,不要基于它推断前置状态;错误消息字符串为固定文案,也不要把error当真实异常对象处理(虽然签名允许None | str | BaseException);
  4. 监听器内部抛出异常会被吞掉并记录日志,不影响状态变更的提交,因此请勿把关键业务逻辑放在监听器异常后的"补救"里;
  5. 本修复针对的是"手动设置 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),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询