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

如何在Airflow的小时调度DAG中设置每日执行任务?

在Airflow中实现小时调度DAG的差异化任务执行频率

完全可以实现。核心思路是基于DAG的小时调度实例,通过逻辑判断控制任务C仅在每日指定时间的调度实例中执行,以下是几种可行方案:

方案一:用ShortCircuitOperator控制任务C的执行

这种方式会直接短路任务C的执行路径,不符合条件时整个下游(这里只有C)都会被跳过:

  1. 定义DAG的调度间隔为@hourly,确保每小时触发一次
  2. 任务A、B保持常规配置,每小时调度触发后正常执行
  3. 在任务B和C之间添加一个ShortCircuitOperator,通过判断当前调度的execution_date决定是否放行任务C

示例代码:

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

def should_run_c(execution_date):
    # 仅在每日0点的调度实例中执行任务C,可根据需求修改小时值
    return execution_date.hour == 0

with DAG(
    dag_id="hourly_dag_daily_c",
    schedule_interval="@hourly",
    start_date=datetime(2024, 1, 1),
    catchup=False  # 关闭补跑,避免历史实例触发不必要的C任务
) as dag:
    task_a = PythonOperator(
        task_id="task_a",
        python_callable=lambda: print("Executing Task A")
    )

    task_b = PythonOperator(
        task_id="task_b",
        python_callable=lambda: print("Executing Task B")
    )

    check_run_condition = ShortCircuitOperator(
        task_id="check_run_c_condition",
        python_callable=should_run_c,
        op_kwargs={"execution_date": "{{ execution_date }}"}
    )

    task_c = PythonOperator(
        task_id="task_c",
        python_callable=lambda: print("Executing Task C")
    )

    task_a >> task_b >> check_run_condition >> task_c

方案二:在任务C的执行逻辑中直接判断

不需要额外的控制任务,直接在任务C的核心逻辑里加入判断,不符合条件时提前退出:

def run_task_c(execution_date):
    if execution_date.hour != 0:
        print("Skip Task C: Not daily execution time")
        return
    # 这里编写任务C的核心业务逻辑
    print("Running Task C core logic")

# 任务C的配置
task_c = PythonOperator(
    task_id="task_c",
    python_callable=run_task_c,
    op_kwargs={"execution_date": "{{ execution_date }}"}
)

方案三:用BranchPythonOperator实现分支执行

如果需要在Airflow UI上清晰区分“执行C”和“跳过C”的状态,可以用分支操作符:

from airflow.operators.dummy import DummyOperator
from airflow.operators.python import BranchPythonOperator

def decide_run_c(execution_date):
    return "task_c" if execution_date.hour == 0 else "skip_task_c"

branch_operator = BranchPythonOperator(
    task_id="branch_run_c",
    python_callable=decide_run_c,
    op_kwargs={"execution_date": "{{ execution_date }}"}
)

skip_task_c = DummyOperator(task_id="skip_task_c")

# 依赖关系配置
task_a >> task_b >> branch_operator >> [task_c, skip_task_c]

注意事项

  • 若不需要补跑历史任务,务必将DAG的catchup参数设为False,避免历史小时调度实例触发任务C
  • 可根据实际需求调整判断逻辑,比如按日期的其他维度(如每月1号)控制任务C的执行频率

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.18 00:45:00