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
相关产品推荐
相关产品推荐

