Apache Airflow循环创建任务XCom获取异常问题解决
解决Airflow循环创建任务时XCom拉取错误的问题
问题现象
在单个Airflow DAG中通过for循环创建多组逻辑相同的任务,预期每组任务(a、b、c)的分支任务能拉取对应前置任务的XCom进行判断,但实际所有分支任务都仅拉取到最后一组(c)任务的XCom结果。
问题根源
这是Python闭包延迟绑定导致的问题:循环中的name变量不会被每个任务立即捕获,当Airflow实际执行任务时,循环已经执行完毕,所有任务都会引用循环最后一次迭代的name值(即'c'),因此所有分支任务都会错误拉取count_task_c的XCom。
解决方案
核心思路是让每个任务持有独立的name标识,而非依赖外部循环变量:
- 给每个Operator的
params参数传入当前循环的name值 - 在任务的可调用函数中,从
context['params']获取对应任务的name - 基于获取到的
name拼接正确的任务ID,拉取对应前置任务的XCom
修改后的完整代码
from airflow import DAG from airflow.operators.python import PythonOperator, BranchPythonOperator from datetime import days_ago default_args = {'start_date': days_ago(1)} dag = DAG( dag_id='batch_test', default_args=default_args, schedule_interval=None ) def count(**context): name = context['params']['name'] count_dict = {'a':50, 'b':100, 'c':150} if count_dict[name] < 100: return f'add_{name}' else: return f'times_{name}' def branch(**context): # 从params中获取当前任务对应的name name = context['params']['name'] # 拉取对应count任务的XCom task_id = context['ti'].xcom_pull(task_ids=f'count_task_{name}') return task_id def add(**context): name = context['params']['name'] # 拉取对应branch任务的XCom ans_key = context['ti'].xcom_pull(task_ids=f'branch_task_{name}') ans_dict = {'add_a':50+100, 'add_b':100+100, 'add_c':150+100} print(ans_dict[ans_key]) def times(**context): name = context['params']['name'] ans_key = context['ti'].xcom_pull(task_ids=f'branch_task_{name}') ans_dict = {'times_a':50*100, 'times_b':100*100, 'times_c':150*100} print(ans_dict[ans_key]) name_list = ['a','b','c'] for name in name_list: exec_count_task = PythonOperator( task_id=f'count_task_{name}', python_callable=count, provide_context=True, params={'name': name}, dag=dag ) exec_branch_task = BranchPythonOperator( task_id=f'branch_task_{name}', python_callable=branch, provide_context=True, params={'name': name}, # 传递当前任务的name dag=dag ) exec_add_count = PythonOperator( task_id=f'add_{name}', python_callable=add, provide_context=True, params={'name': name}, # 传递当前任务的name dag=dag ) exec_times_count = PythonOperator( task_id=f'times_{name}', python_callable=times, provide_context=True, params={'name': name}, # 传递当前任务的name dag=dag ) exec_count_task >> exec_branch_task >> [exec_add_count, exec_times_count]
关键修改点
- 给
BranchPythonOperator、add和times对应的PythonOperator都添加了params={'name': name},确保每个任务持有独立的标识 - 在
branch、add、times函数中,从context['params']['name']获取当前任务的name,避免依赖外部循环变量 - 拉取XCom时,用获取到的
name拼接正确的前置任务ID,保证每组任务拉取对应前置的XCom结果
内容的提问来源于stack exchange,提问作者DG A
相关产品推荐
相关产品推荐

