You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Airflow手动设置任务/DAG状态时如何触发回调?

Airflow手动修改任务/DAG状态触发回调的解决方案

默认的on_failure_callback和on_success_callback确实仅在任务由调度器自动执行完成时触发,手动修改状态不会触发这些回调——因为Airflow的核心回调逻辑绑定在任务的执行生命周期流程内,手动状态变更属于外部操作,不在该流程覆盖范围内。

以下是几种可行的实现方式:

1. 使用Airflow事件监听机制(推荐,Airflow 2.3+支持)

Airflow提供的EventListener接口可以监听任务状态变更事件,无论状态是调度器自动设置还是手动修改,都能被捕获。

实现步骤:

  • 自定义事件监听器类,继承BaseEventListener并重写状态变更处理方法:
from airflow.listeners import BaseEventListener
from airflow.models import TaskInstance
from airflow.utils.state import State

class ManualTaskStateListener(BaseEventListener):
    def on_task_instance_state_change(
        self,
        previous_state: State | None,
        new_state: State,
        task_instance: TaskInstance,
    ) -> None:
        # 通过run_id区分手动操作(手动修改的任务run_id通常包含"__manual__"标识)
        if "__manual__" in task_instance.run_id:
            # 执行自定义回调逻辑,比如告警
            if new_state == State.FAILED:
                self.trigger_alert(task_instance, "failed")
            elif new_state == State.SUCCESS:
                self.trigger_alert(task_instance, "success")

    def trigger_alert(self, ti: TaskInstance, state: str):
        # 此处编写告警逻辑,例如发送邮件、企业微信通知等
        print(f"手动修改任务状态:DAG [{ti.dag_id}] 任务 [{ti.task_id}] 被设置为{state}")
  • 在Airflow配置文件airflow.cfg中注册监听器:
listener_classes = 你的模块路径.ManualTaskStateListener

2. 手动调用API时附加回调逻辑

如果你的状态修改是通过Airflow REST API完成的,可以在调用修改状态的接口后,主动触发回调代码。例如调用POST /taskInstances/{dag_id}/{task_id}/{execution_date}/state接口后,立即执行告警或回调逻辑。

3. 数据库触发器(适合复杂场景)

Airflow的任务状态存储在task_instance表中,可以通过数据库触发器(如PostgreSQL的TRIGGER)监听该表state字段的变更。当检测到状态变更且符合手动操作特征(如run_id包含__manual__)时,调用外部脚本执行回调逻辑。

注意事项:

  • 手动修改的任务记录,run_id通常带有__manual__标识,可用于区分自动执行与手动操作;
  • 事件监听机制仅支持Airflow 2.3及以上版本,低版本需使用其他方案;
  • 自定义回调逻辑需保证幂等性,避免重复触发告警。

内容的提问来源于stack exchange,提问作者Ares

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 01:19:51