Airflow中sensor子任务能否触发调度器重新运行自身及提交报告的父任务?
实现方案
完全可以实现,不需要触发任务失败即可完成父任务和当前sensor任务的状态重置与重跑,Airflow原生提供了公开API支持这类操作,以下是具体实现方式:
核心实现逻辑
在sensor检测到业务系统报告丢失的分支逻辑中,直接修改对应任务实例的状态为待调度状态,Airflow scheduler会自动识别状态变更并重新调度这两个任务执行。
具体操作步骤
- 首先在sensor的执行逻辑中获取当前DAG运行实例(DagRun),以及需要重置的两个任务的实例对象:提交报告参数的父任务、当前sensor任务
- 将两个任务实例的状态设置为
None(即清空运行记录,标记为待调度) - 无需额外给scheduler发自定义信号,scheduler会自动扫描到状态变更的任务,加入执行队列
代码示例
你可以直接在sensor的poke逻辑中加入以下处理代码:
from airflow.utils.state import State from airflow.operators.python import get_current_context # 以下为sensor的poke方法内的逻辑 def poke(self, context): # 原有检查报告是否生成的业务逻辑 report_available = check_business_system_report() if not report_available: # 识别到报告丢失,重置目标任务状态 dag_run = context["dag_run"] # 替换为你实际的提交报告参数的父任务ID parent_task_id = "submit_report_params" # 重置父任务状态 parent_ti = dag_run.get_task_instance(parent_task_id) parent_ti.set_state(State.NONE) # 重置当前sensor任务状态 current_ti = context["ti"] current_ti.set_state(State.NONE) # 返回False即可,scheduler会自动触发重跑 return False # 报告正常存在的逻辑 return True
注意事项
- sensor建议配置为
reschedule模式,避免长期占用worker资源,重跑调度效率更高 - 建议增加最大重跑次数限制,避免业务系统长期异常时任务无限重跑,超过阈值后再标记为失败触发告警
- 所有操作均使用Airflow公开的TaskInstance、DagRun API,无需调用scheduler内部私有接口,版本兼容性更强
内容的提问来源于stack exchange,提问作者Tevett Goad
相关产品推荐
相关产品推荐

