如何在Airflow中让任务暂停24小时后重复执行?
实现Airflow DAG中task1→等待24小时→task1的方案
当然可以实现这个需求,下面提供两种可行的方案,优先推荐第一种原生传感器的方式:
方案一:使用TimeDeltaSensor(推荐)
Airflow内置的TimeDeltaSensor是专门用于等待指定时间间隔的组件,不会长时间占用Worker资源,适合这种24小时的长等待场景。
代码示例
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.sensors.time_delta import TimeDeltaSensor from datetime import datetime, timedelta # 定义API拉取任务的逻辑 def pull_api_data(): # 这里替换为你的实际API拉取代码 print("完成API数据拉取") default_args = { 'owner': 'airflow', 'start_date': datetime(2024, 1, 1), 'retries': 0, } with DAG( dag_id='api_daily_pull_with_wait', default_args=default_args, schedule_interval=None, # 手动触发,若需定期启动可改为@daily等 catchup=False, ) as dag: # 第一次执行API拉取 first_pull = PythonOperator( task_id='first_api_pull', python_callable=pull_api_data ) # 等待24小时的传感器任务 wait_24h = TimeDeltaSensor( task_id='wait_24_hours', delta=timedelta(hours=24), mode='reschedule' # reschedule模式会释放Worker资源,定期检查时间是否到点 ) # 第二次执行API拉取 second_pull = PythonOperator( task_id='second_api_pull', python_callable=pull_api_data ) # 设置任务依赖 first_pull >> wait_24h >> second_pull
关键说明
mode='reschedule':相比默认的poke模式,会在等待期间释放Worker资源,每隔一段时间(默认30秒)重新调度检查,更适合长时间等待场景。schedule_interval=None:如果只需要手动触发一次执行两次拉取,保持这个设置;如果需要每天自动启动一次这个流程,可改为@daily。
方案二:PythonOperator内嵌sleep(不推荐)
这种方式直接在Python任务中用time.sleep()实现等待,但会持续占用Worker进程24小时,浪费资源,仅适合测试或资源充足的场景。
代码示例
from airflow import DAG from airflow.operators.python import PythonOperator from datetime import datetime import time def pull_api_data(): print("完成API数据拉取") def wait_24h(): time.sleep(86400) # 24小时对应的秒数 default_args = { 'owner': 'airflow', 'start_date': datetime(2024, 1, 1), 'retries': 0, } with DAG( dag_id='api_pull_with_sleep', default_args=default_args, schedule_interval=None, catchup=False, ) as dag: first_pull = PythonOperator(task_id='first_pull', python_callable=pull_api_data) wait_task = PythonOperator(task_id='wait_24h', python_callable=wait_24h) second_pull = PythonOperator(task_id='second_pull', python_callable=pull_api_data) first_pull >> wait_task >> second_pull
注意事项
- 如果Airflow在等待期间重启,这个sleep任务会中断,需要重新执行;而
TimeDeltaSensor会基于元数据库记录的时间继续等待,可靠性更高。
内容的提问来源于stack exchange,提问作者Nutzuchima
相关产品推荐
相关产品推荐

