Airflow无法获取Task Instance读取XCom的问题求助
问题根源
你遇到的KeyError: 'ti'是因为在PythonOperator中使用lambda: do_task()的写法阻断了Airflow上下文参数的传递。Airflow确实会自动将Task Instance(ti)等上下文参数注入到python_callable的kwargs中,但你的lambda表达式没有把这些参数传递给do_task函数,导致do_task内部的kwargs为空,自然找不到'ti'键。
解决方案
直接将python_callable指定为do_task即可,不需要用lambda包装,Airflow会自动把上下文参数和op_kwargs里的参数一起传递给do_task函数。
修改后的task_1和task_2代码如下:
task_1 = PythonOperator( task_id="task_1", do_xcom_push=False, op_kwargs={ 'connector_key':"e_conns", }, python_callable=do_task, # 移除lambda,直接指向函数 ) task_2 = PythonOperator( task_id="task_2", do_xcom_push=False, op_kwargs={ 'connector_key':"f_conns", }, python_callable=do_task, # 同样修改此处 )
如果确实需要用lambda(非必要场景),必须显式传递kwargs:
python_callable=lambda **kwargs: do_task(**kwargs),
补充说明
你的do_task函数已经定义为def do_task(**kwargs) -> int:,只要保证PythonOperator直接调用该函数(或通过lambda传递参数),Airflow就会自动将ti、dag_run等上下文参数注入到kwargs中,你就能正常通过kwargs['ti']获取Task Instance来拉取XCom。
内容的提问来源于stack exchange,提问作者FrustratedWithFormsDesigner
相关产品推荐
相关产品推荐

