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

Airflow中PythonSensor超时后是否可以触发其他指定任务?

PythonSensor超时触发指定任务实现方案

完全可以实现该需求,你之前使用BranchPythonOperator的方案冗余且会破坏原有重试特性,推荐采用trigger_rule配置的方式实现,同时保留传感器的自动重试能力,具体实现逻辑如下:

  • 给PythonSensor单独配置参数:设置timeout为你需要的单次监听超时时间,同时配置retries(重试次数)、retry_delay(重试间隔)参数,保留首次失败自动重调度的能力
  • 正常业务任务保持默认trigger_rule="all_success",依赖设置为该PythonSensor,传感器监听成功后自动触发
  • 超时要触发的指定任务设置trigger_rule="all_failed",依赖同样设置为该PythonSensor,传感器用尽重试次数仍超时失败时,会自动触发该任务,传感器监听成功时该任务会自动跳过

示例代码片段

from airflow import DAG
from airflow.sensors.python import PythonSensor
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta

# 你的FTP监听逻辑
def check_ftp_file():
    # 这里写FTP服务器文件存在性检查逻辑,存在返回True,不存在返回False
    pass

# 超时触发的任务逻辑
def handle_timeout():
    # 这里写超时后要执行的指定任务逻辑
    pass

# 正常业务任务逻辑
def normal_business():
    # 这里写文件存在后要执行的正常业务逻辑
    pass

with DAG(
    dag_id="ftp_monitor_dag",
    start_date=datetime(2024, 1, 1),
    schedule_interval="@daily",
    catchup=False
) as dag:
    ftp_sensor = PythonSensor(
        task_id="monitor_ftp_file",
        python_callable=check_ftp_file,
        timeout=3600,  # 单次监听超时时间1小时
        retries=2,  # 失败后重试2次
        retry_delay=timedelta(minutes=10),  # 每次重试间隔10分钟
        poke_interval=60,  # 每60秒检查一次FTP
        mode="poke"
    )

    timeout_task = PythonOperator(
        task_id="handle_sensor_timeout",
        python_callable=handle_timeout,
        trigger_rule="all_failed"  # 上游传感器全失败时才触发
    )

    business_task = PythonOperator(
        task_id="run_normal_business",
        python_callable=normal_business
    )

    ftp_sensor >> [timeout_task, business_task]

以上方案完全满足需求,不需要额外引入分支算子,传感器本身的重试逻辑可以正常生效。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.24 23:45:02