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

