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

Airflow如何实现任务运行15分钟后再执行下游任务?

原有方案失效原因

你写的time.sleep(900)逻辑从根上就不符合Airflow的调度规则:你配置的DAG调度间隔是1分钟,意味着每分钟都会生成一个独立的DAG Run实例,每个实例里的休眠任务都会单独卡15分钟,既做不到累计15次写入再触发下游,还会长时间占用Worker进程槽位,大概率触发任务运行超时、Worker资源占满的问题,完全达不到预期效果。

可行实现方案

不要在任务代码里写硬休眠,用Airflow原生能力实现,两种方案按需选:

方案1:双DAG+数据集调度(Airflow 2.4及以上版本首选,生产环境推荐)

把写数据和分析数据拆成两个独立DAG,用Airflow的数据集能力做触发联动:

  • 第一个DAG保持1分钟调度,只负责每分钟执行写文件逻辑,声明产出对应数据文件的数据集
  • 第二个DAG绑定该数据集作为触发条件,配置15分钟的累计窗口,等数据集攒够15分钟的写入更新后,自动触发后续分析任务
    参考代码:
from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.datasets import Dataset
from datetime import datetime, timedelta

# 定义采集数据文件对应的数据集标识
COLLECTED_DATA = Dataset("file:///path/to/your/collected_data.log")

# DAG1:每分钟执行一次数据写入
with DAG(
    dag_id="collect_data_per_min",
    schedule_interval="* * * * *",
    start_date=datetime(2024, 1, 1),
    catchup=False
) as collect_dag:
    write_task = PythonOperator(
        task_id="write_data_to_file",
        python_callable=your_existing_write_function, # 替换成你实际写文件的业务函数
        outlets=[COLLECTED_DATA]
    )

# DAG2:累计15分钟写入后自动触发分析
with DAG(
    dag_id="analyze_data_per_15min",
    schedule=COLLECTED_DATA,
    start_date=datetime(2024, 1, 1),
    catchup=False,
    dataset_condition_timedelta=timedelta(minutes=15)
) as analyze_dag:
    downstream_task = PythonOperator(
        task_id="run_analysis",
        python_callable=your_existing_downstream_function # 替换成你的后续业务任务
    )

方案2:单DAG短路判断(适配2.4以下低版本Airflow)

如果版本不支持数据集特性,保持单DAG1分钟调度,在下游任务前加短路判断逻辑:每次DAG运行时先检查是否已经累计采集满15分钟,不满就直接跳过下游任务,满了才放行执行分析,同时重置计数进入下一轮采集周期。
参考代码:

import os
import time
from airflow import DAG
from airflow.operators.python import PythonOperator, ShortCircuitOperator
from datetime import datetime

TIMESTAMP_MARK_PATH = "/tmp/collect_start_ts.mark"

def write_data_with_mark():
    # 第一次写入时记录采集起始时间
    if not os.path.exists(TIMESTAMP_MARK_PATH):
        with open(TIMESTAMP_MARK_PATH, "w") as f:
            f.write(str(time.time()))
    # 执行原有的每分钟写文件逻辑
    your_existing_write_function()

def check_collect_duration():
    if not os.path.exists(TIMESTAMP_MARK_PATH):
        return False
    with open(TIMESTAMP_MARK_PATH, "r") as f:
        start_ts = float(f.read())
    # 累计满15分钟就放行,同时删除标记重置下一轮
    if time.time() - start_ts >= 900:
        os.remove(TIMESTAMP_MARK_PATH)
        return True
    return False

with DAG(
    dag_id="collect_and_analyze",
    schedule_interval="* * * * *",
    start_date=datetime(2024, 1, 1),
    catchup=False
) as my_dag:
    task1 = PythonOperator(
        task_id="write_data_per_min",
        python_callable=write_data_with_mark
    )
    check_task = ShortCircuitOperator(
        task_id="check_if_15min_collected",
        python_callable=check_collect_duration
    )
    task2 = PythonOperator(
        task_id="run_downstream_tasks",
        python_callable=your_existing_downstream_function
    )

    task1 >> check_task >> task2
避坑提醒

不要在Airflow任务里写长耗时sleep逻辑:Airflow Worker的并发槽位是有限的,长睡眠任务会一直占用槽位不释放,很容易把Worker堵死,导致其他任务无法调度。如果业务逻辑允许,也可以直接把DAG调度间隔改成15分钟,在写数据任务内循环15次每分钟完成一次写入再跑下游,但这种方案容错性差,写入中途出错会丢失整轮采集数据,不建议生产环境使用。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.27 19:39:14