Airflow中Python方法访问XCom报ti/task_instance KeyError如何解决
问题根因
报错核心原因是python_callable参数传值错误:你当前写法是refresh_task(task_id, partitions, dag, iteration),等于在DAG文件解析阶段就直接执行了refresh_task函数,而非将函数引用传递给PythonOperator等待任务调度运行时调用。DAG解析阶段不存在任务运行时的上下文对象,自然无法从**context中取到task_instance或ti,触发KeyError。
除此之外代码还有几处隐性问题:
- 循环内未定义
iteration变量,直接使用会触发未定义报错 foo是refresh_task内部定义的局部变量,写task_group.foo属于错误的属性引用,运行时会抛AttributeError- Airflow 2.0+版本已自动完成上下文注入,不需要手动设置
provide_context=True
修复方案
- 首先导入
functools.partial用于提前绑定函数的固定参数,避免在DAG解析阶段直接执行调用函数:
from functools import partial
- 修正任务组循环内的PythonOperator定义,使用enumerate遍历partition_list拿到迭代序号,将绑定完固定参数的函数引用传给python_callable:
# 用enumerate拿到循环序号iteration for iteration, partitions in enumerate(partition_list): t1 = PythonOperator( task_id=f'Refresh_Wrapper_{iteration}', python_callable=partial( refresh_task, task_group=task_id, data=partitions, dag=dag, iteration=iteration ), dag=dag )
- 修正refresh_task内部的逻辑,修正变量引用错误,正确从上下文获取任务实例拉取XCom:
def refresh_task(task_group, data, dag, iteration, **context): # 从上下文获取任务实例对象,ti和task_instance两个键完全等价 ti = context['ti'] foo = f"baz{iteration}" # foo是当前函数内的局部变量,直接传入xcom_pull即可,不需要加task_group前缀 pulled_xcom_value = ti.xcom_pull(task_ids=foo, key='return_value') bar = f"Bar {pulled_xcom_value}" # 其余业务逻辑按原有需求编写
排查补充
如果修改后仍然无法拉取到XCom值,检查两项配置:
- 被拉取XCom的上游任务(task_id为
baz{iteration})必须和当前任务设置了正确的依赖关系,确保上游任务先于当前任务执行 - 上游任务必须通过return返回值、或主动调用
ti.xcom_push()的方式写入了对应key的XCom数据
内容的提问来源于stack exchange,提问作者Deslyxia
相关产品推荐
相关产品推荐

