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

如何基于调度规则为Airflow DAG配置不同参数进行调度

如何在Airflow中根据调度日期动态设置DAG参数?

这个需求我之前在项目里也碰到过,核心思路就是基于Airflow的执行日期(execution_date)来做分支判断,根据不同的星期几给source变量赋值。下面给你两种实用的实现方案,你可以根据自己的场景选:

方案一:在单个Operator内动态判断(适合单一任务使用的场景)

如果只有一个任务需要用到source变量,直接在Operator的业务逻辑里判断执行日期的星期几就行。Airflow自带的pendulum库处理日期非常方便,它的weekday()方法返回0代表周一,6代表周日,刚好对应我们的需求:

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta
import pendulum

def process_data(**context):
    # 从上下文获取执行日期并转成pendulum对象
    exec_date = pendulum.parse(context['execution_date'].isoformat())
    weekday = exec_date.weekday()
    
    # 根据星期几赋值source,注意优先级:周二同时出现在两个规则里,这里先判断周一/周二,所以周二会走process_usa
    if weekday in [0, 1]:  # 周一、周二
        source = "process_usa"
    elif weekday == 4:  # 周五
        source = "process_ca"
    elif weekday in [1, 2, 6]:  # 周二、周三、周日
        source = "process_ind"
    else:
        # 其他日期可以设置默认值或者直接终止任务
        source = None
        print("当前日期没有匹配的source规则,任务终止")
        return
    
    # 这里写你的实际业务逻辑,比如调用处理脚本、读写数据库等
    print(f"开始执行任务,source参数为:{source}")

# DAG基础配置
default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2024, 1, 1),
    'retries': 1,
    'retry_delay': timedelta(minutes=5),
}

with DAG(
    'dynamic_source_single_task',
    default_args=default_args,
    schedule_interval='@daily',  # 每天调度,由代码内部过滤日期规则
    catchup=False,
) as dag:
    run_process = PythonOperator(
        task_id='run_process_task',
        python_callable=process_data,
        provide_context=True,  # 必须开启,才能获取execution_date等上下文变量
    )

run_process

如果是用BashOperator,也可以通过Jinja模板直接在bash命令里做判断:

{% set exec_date = execution_date.in_timezone('UTC') %}
{% set weekday = exec_date.weekday() %}
{% if weekday in [0,1] %}
export SOURCE=process_usa
{% elif weekday ==4 %}
export SOURCE=process_ca
{% elif weekday in [1,2,6] %}
export SOURCE=process_ind
{% endif %}

# 执行你的业务脚本
./your_business_script.sh $SOURCE

方案二:全局设置参数(适合多个任务共用source的场景)

如果DAG里有多个任务都需要用到source变量,推荐先通过一个前置任务计算好source,然后用XCom传递给后续任务,这样所有任务都能复用这个参数:

from airflow import DAG
from airflow.operators.python import PythonOperator
from datetime import datetime, timedelta
import pendulum

def calculate_source(**context):
    exec_date = pendulum.parse(context['execution_date'].isoformat())
    weekday = exec_date.weekday()
    
    # 同样注意优先级,这里调整了顺序,周二会优先走process_ind
    if weekday == 4:
        source = "process_ca"
    elif weekday in [1, 2, 6]:
        source = "process_ind"
    elif weekday in [0, 1]:
        source = "process_usa"
    else:
        source = None
    
    # 将source推送到XCom,供后续任务获取
    context['ti'].xcom_push(key='task_source', value=source)

def task_one(**context):
    # 从XCom拉取source参数
    source = context['ti'].xcom_pull(key='task_source', task_ids='calculate_source_task')
    if not source:
        print("未获取到有效source参数,任务终止")
        return
    print(f"任务1开始执行,source: {source}")

def task_two(**context):
    source = context['ti'].xcom_pull(key='task_source', task_ids='calculate_source_task')
    print(f"任务2开始执行,source: {source}")

default_args = {
    'owner': 'airflow',
    'depends_on_past': False,
    'start_date': datetime(2024, 1, 1),
    'retries': 1,
}

with DAG(
    'dynamic_source_global',
    default_args=default_args,
    schedule_interval='@daily',
    catchup=False,
) as dag:
    # 前置任务:计算并推送source
    calc_source = PythonOperator(
        task_id='calculate_source_task',
        python_callable=calculate_source,
        provide_context=True,
    )
    
    # 后续任务:获取source并执行
    task_1 = PythonOperator(
        task_id='task_one',
        python_callable=task_one,
        provide_context=True,
    )
    
    task_2 = PythonOperator(
        task_id='task_two',
        python_callable=task_two,
        provide_context=True,
    )

# 设置任务依赖
calc_source >> [task_1, task_2]

重要提示

注意你需求里周二同时属于两个规则,一定要明确优先级!比如上面的代码里,哪个判断逻辑写在前面,周二就会走哪个分支。如果需要周二同时执行两个不同source的任务,那就要用BranchPythonOperator来分支任务,而不是单纯设置变量。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.30 07:53:11