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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.20 18:32:19