Airflow 2.5.1中如何传递PythonOperator任务的返回值
在Airflow 2.5.1中传递上游任务返回值给下游任务
Airflow的PythonOperator默认会把任务函数的返回值以return_value为key推送到XCom中,你只需要通过XCom拉取并传递给下游任务即可。同时必须明确任务间的依赖关系,确保下游任务在上游任务执行完成后再运行。下面提供两种可行的实现方式:
方法一:在任务函数中通过Task Instance拉取XCom
修改下游任务的函数,接收ti(Task Instance)参数,直接从XCom中拉取上游任务的返回值:
from airflow.operators.python import PythonOperator from airflow.models import DAG from datetime import datetime, timedelta def add(): return 1 + 1 def multiply(ti): # 拉取t1任务的XCom返回值,默认key为return_value a = ti.xcom_pull(task_ids='t1') return a * 999 dag_args = { 'owner': 'me', 'depends_on_past': False, 'start_date': datetime(2023, 2, 27), 'email': ['me@myhome.com'], 'email_on_failure': True, 'email_on_retry': True, 'retries': 1, 'retry_delay': timedelta(minutes=3)} with DAG( dag_id='dag', start_date=datetime(2023, 2, 27), default_args=dag_args, schedule_interval='@once', end_date=None,) as dag: t1 = PythonOperator( task_id="t1", python_callable=add, dag=dag ) t2 = PythonOperator( task_id="t2", python_callable=multiply, # 开启上下文传递,让函数能接收ti参数 provide_context=True, dag=dag ) # 设置任务依赖,确保t2在t1执行完成后运行 t1 >> t2
方法二:通过模板语法直接传递参数到下游任务
不需要修改任务函数,而是在定义下游任务时,用Jinja2模板语法将上游XCom值传入op_kwargs:
from airflow.operators.python import PythonOperator from airflow.models import DAG from datetime import datetime, timedelta def add(): return 1 + 1 def multiply(a): return a * 999 dag_args = { 'owner': 'me', 'depends_on_past': False, 'start_date': datetime(2023, 2, 27), 'email': ['me@myhome.com'], 'email_on_failure': True, 'email_on_retry': True, 'retries': 1, 'retry_delay': timedelta(minutes=3)} with DAG( dag_id='dag', start_date=datetime(2023, 2, 27), default_args=dag_args, schedule_interval='@once', end_date=None,) as dag: t1 = PythonOperator( task_id="t1", python_callable=add, dag=dag ) t2 = PythonOperator( task_id="t2", python_callable=multiply, # 用模板语法拉取t1的XCom值,作为a参数传入 op_kwargs={'a': '{{ ti.xcom_pull(task_ids="t1") }}'}, dag=dag ) t1 >> t2
注意事项
- 必须设置
t1 >> t2的依赖关系,否则下游任务可能在上游任务未执行时就启动,导致拉取不到XCom值。 - PythonOperator默认开启
do_xcom_push=True,会自动将函数返回值推送到XCom;如果需要自定义XCom的key,可以在任务函数中用ti.xcom_push(key='自定义key', value=值)主动推送,拉取时指定对应的key即可。
内容的提问来源于stack exchange,提问作者FxxkDogeCoins
相关产品推荐
相关产品推荐

