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

如何设置Airflow的schedule_interval使其仅在每小时0-30分钟且前序任务结束后运行?

解决Airflow DAG仅在每小时0-30分钟时段内调度的问题

你之前用的0-30/60 * * * * Cron表达式逻辑有误——它表示每小时的0-30分钟区间内,每60分钟执行一次,实际只会在整点触发,完全不符合你要的“在0-30分钟时段内随时触发(且需等前一次调度结束)”的需求。下面给两种实用的实现方案:

方案一:分支任务+时间/前置运行检查

通过高频调度触发DAG,再用分支任务判断两个核心条件:当前时间是否在每小时0-30分钟内、前一次DAG运行是否成功,满足才执行主业务逻辑。

代码示例

from airflow import DAG
from airflow.operators.python import PythonOperator, BranchPythonOperator
from airflow.utils.dates import days_ago
from airflow.models import DagRun
from datetime import datetime

def check_runtime_conditions(**context):
    # 检查当前时间是否在0-30分钟区间
    current_minute = datetime.now().minute
    if current_minute not in range(0, 31):
        return 'skip_execution'
    
    # 检查前一次DAG运行状态(首次运行跳过)
    dag_id = context['dag'].dag_id
    previous_runs = DagRun.find(dag_id=dag_id, execution_date__lt=context['execution_date'])
    if previous_runs:
        last_successful_run = max(previous_runs, key=lambda x: x.execution_date)
        if last_successful_run.state != 'success':
            return 'skip_execution'
    
    # 条件全部满足,执行主任务
    return 'execute_main_task'

def skip_task():
    print("当前不在允许时段或前一次运行未完成,跳过执行")

def main_business_task():
    # 替换成你的实际业务逻辑
    print("执行核心业务任务")

with DAG(
    dag_id='hourly_restricted_first_30min',
    schedule_interval='*/5 * * * *',  # 每5分钟触发一次,可根据需求调整频率
    start_date=days_ago(1),
    catchup=False,  # 关闭历史补跑
    tags=['hourly', 'time_restricted']
) as dag:

    condition_check = BranchPythonOperator(
        task_id='check_runtime_conditions',
        python_callable=check_runtime_conditions,
        provide_context=True
    )

    skip = PythonOperator(
        task_id='skip_execution',
        python_callable=skip_task
    )

    main_task = PythonOperator(
        task_id='execute_main_task',
        python_callable=main_business_task
    )

    condition_check >> [main_task, skip]

方案二:依赖前置运行+时间窗口校验

利用Airflow原生的depends_on_past属性确保前一次运行成功才触发,再结合时间校验任务过滤非允许时段的调度。

代码示例

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.operators.latest_only import LatestOnlyOperator
from airflow.utils.dates import days_ago
from datetime import datetime

def validate_time_window(**context):
    current_minute = datetime.now().minute
    if not (0 <= current_minute <= 30):
        raise ValueError("当前不在每小时0-30分钟时段,终止执行")
    print("进入允许执行的时间窗口")

def main_business_task():
    # 替换成你的实际业务逻辑
    print("执行核心业务任务")

with DAG(
    dag_id='hourly_restricted_to_first_30min',
    schedule_interval='*/5 * * * *',
    start_date=days_ago(1),
    catchup=False,
    depends_on_past=True,  # 强制依赖前一次DAG运行成功
    tags=['hourly', 'time_window']
) as dag:

    latest_only = LatestOnlyOperator(task_id='latest_only')  # 仅在最新调度窗口执行

    time_validation = PythonOperator(
        task_id='validate_time_window',
        python_callable=validate_time_window,
        provide_context=True
    )

    main_task = PythonOperator(
        task_id='main_business_task',
        python_callable=main_business_task
    )

    latest_only >> time_validation >> main_task

关键细节说明

  1. 调度频率:设置为*/5 * * * *(每5分钟)是为了在0-30分钟内及时触发,你可以根据需求调整为* * * * *(每分钟)提升灵敏度。
  2. depends_on_past=True:确保只有前一次DAG运行成功,才会触发当前调度,完全满足“前一个调度结束后才运行”的要求。
  3. 时间校验:通过Python代码直接判断当前分钟数,逻辑清晰且灵活,比Cron表达式更易控制。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.17 02:33:24