如何将函数返回的变量传递给DAG及依赖DAG的任务
实现日期参数跨任务、跨DAG传递方案
核心思路
利用Airflow的XCom实现同DAG内任务间的参数传递,再通过TriggerDagRunOperator的conf参数实现跨DAG的参数传递。下面是具体实现步骤:
1. 补全基础依赖与日期生成任务
先完善日期生成函数,并用PythonOperator创建任务生成日期——该算子默认会把函数返回值推送到XCom(默认key为return_value):
from airflow import DAG from airflow.operators.python import PythonOperator from airflow.operators.trigger_dagrun import TriggerDagRunOperator from datetime import datetime, timedelta # 根据实际需求定义要减去的天数 days_to_subtract = 1 def date_fn(): d = datetime.today() - timedelta(days=days_to_subtract) # 转成字符串避免datetime对象序列化问题 return d.strftime("%Y-%m-%d")
2. 修改processthe_files DAG,实现同DAG内参数传递
在processthe_files DAG中添加日期生成任务,让内部任务拉取XCom的日期参数,同时配置触发跨DAG的参数传递:
SCHEDULE = "0 8 * * 4" with DAG( dag_id="processthe_files", start_date=datetime(2024, 10, 8), schedule_interval=SCHEDULE, catchup=False ) as dag: # 1. 生成日期并推送到XCom generate_date = PythonOperator( task_id="generate_date", python_callable=date_fn ) # 2. file_processing任务拉取XCom的日期参数 # 假设你的Job类支持通过parameters传递参数,根据实际情况调整 file_processing = Job( parameters={"target_date": "{{ ti.xcom_pull(task_ids='generate_date') }}"} ).to_task # 3. 触发processtables DAG时,通过conf传递日期 trigger_processtables = TriggerDagRunOperator( task_id='trigger_processtables', trigger_dag_id='processtables', # 和目标DAG的dag_id保持一致 wait_for_completion=True, # 将XCom的日期放入conf,传递给目标DAG conf={"target_date": "{{ ti.xcom_pull(task_ids='generate_date') }}"}, dag=dag ) # 设置任务依赖链 generate_date >> file_processing >> trigger_processtables
3. 修改processtables DAG,接收跨DAG传递的参数
在processtables DAG的任务中,从dag_run.conf中取出传递过来的日期参数:
with DAG( dag_id="processtables", # 和TriggerDagRunOperator中的trigger_dag_id一致 start_date=datetime(2024, 10, 8), schedule_interval=None, catchup=False ) as dag: processtables = Job( # 从dag_run.conf中获取日期参数 parameters={"target_date": "{{ dag_run.conf.get('target_date') }}"} ).to_task processtables
关键细节说明
- XCom自动推送:
PythonOperator默认会把函数返回值推送到XCom,无需手动调用ti.xcom_push(),通过task_ids即可指定拉取的任务。 - 日期序列化:将datetime对象转为字符串传递,避免Airflow序列化datetime对象时出现兼容性问题。
- conf参数传递:
TriggerDagRunOperator的conf会作为触发目标DAG的配置,目标DAG可通过dag_run.conf访问这些参数。 - 原代码修正:你原代码中
processthefiles >> trigger_processtables是变量名错误,已修正为file_processing >> trigger_processtables。
内容的提问来源于stack exchange,提问作者Aviator
相关产品推荐
相关产品推荐

