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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 06:36:05