You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.08.30 19:18:17