如何基于调度规则为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
相关产品推荐
相关产品推荐

