如何通过Airflow UI判断任务已清除及触发对应逻辑?
关于Airflow任务清除相关问题的解答
1. 如何通过Airflow UI查看已被清除的任务?
Airflow UI支持直接筛选查看已清除的任务实例:
- 进入目标DAG的详情页,切换到Task Instances标签页,在筛选栏的State选项中选择
cleared,即可看到该DAG下所有被清除的任务实例; - 若需全局查看所有已清除任务,可点击顶部导航栏的Browse → Task Instances,同样通过
cleared状态筛选。
2. 如何在Airflow UI清除任务时触发特定逻辑?
Airflow 2.x及以上版本可以通过**事件监听器(Listeners)**实现,监听任务清除事件并执行自定义逻辑:
- 编写自定义监听器代码,监听
TaskInstanceClearedEvent事件:
from airflow.listeners import hookimpl from airflow.events import TaskInstanceClearedEvent @hookimpl def on_task_instance_cleared(event: TaskInstanceClearedEvent): # 在此编写你的特定逻辑,例如发送告警、记录审计日志、调用外部服务等 dag_id = event.task_instance.dag_id task_id = event.task_instance.task_id print(f"任务 {task_id}(所属DAG:{dag_id})于 {event.timestamp} 被清除")
- 将上述代码保存为Python文件,放置到Airflow的
plugins目录下,Airflow会自动加载该监听器,之后每次通过UI清除任务时,都会触发你定义的逻辑。
3. 任务运行时如何判断是因UI清除而重启?
可以结合XCom标记或任务实例状态历史来实现:
方法1:通过XCom标记(推荐)
配合上述的任务清除监听器,在任务被清除时写入XCom标记:
# 在监听器的on_task_instance_cleared函数中添加 event.task_instance.xcom_push(key='is_cleared_restart', value=True)
然后在任务代码中读取该标记:
def my_task(**context): ti = context['ti'] is_cleared_restart = ti.xcom_pull(key='is_cleared_restart', task_ids=ti.task_id) if is_cleared_restart: # 处理清除后重启的逻辑 print("当前任务是因被清除而重启")
方法2:检查任务实例历史状态
在任务代码中查询当前任务实例的上一次尝试记录,判断其状态是否为cleared:
from airflow.models import TaskInstance from airflow.utils.session import create_session def check_cleared_restart(**context): ti = context['ti'] with create_session() as session: prev_ti = session.query(TaskInstance).filter( TaskInstance.dag_id == ti.dag_id, TaskInstance.task_id == ti.task_id, TaskInstance.execution_date == ti.execution_date, TaskInstance.try_number == ti.try_number - 1 ).first() return prev_ti is not None and prev_ti.state == 'cleared' def my_task(**context): if check_cleared_restart(**context): print("当前任务是因被清除而重启")
内容的提问来源于stack exchange,提问作者Eugene Biruk
相关产品推荐
相关产品推荐

