Airflow获取Python Operator前置任务返回值报错(KeyError: 'ti')求助
Airflow中Python Operator获取前置任务返回值及KeyError: 'ti'问题解决
问题分析
你的代码存在两个核心问题导致报错及功能失效:
- PythonOperator循环引用:定义的
pythondefination变量和它指定的python_callable同名,变量未完成定义就被引用,会引发未定义错误,后续逻辑也无法正常执行。 - 未为
valuecapture创建对应Operator且未传递上下文:valuecapture仅被定义但未关联到任何PythonOperator任务,即使关联,也未配置上下文传递参数,导致无法获取ti(任务实例对象)。
修复后的完整代码示例
from datetime import datetime from airflow import DAG from airflow.operators.python import PythonOperator # 前置任务的业务函数,与Operator变量名区分开 def python_defination_func(**kwargs): # 替换为你的实际业务逻辑,返回值会自动推送到XCom return {"current_time": datetime.now().strftime("%Y-%m-%d %H:%M:%S")} def value_capture(**kwargs): ti = kwargs['ti'] # 拉取前置任务的XCom返回值 time_info = ti.xcom_pull(task_ids='python_defination_task') print(f"前置任务返回值: {time_info}") return time_info with DAG( dag_id="test_dag", start_date=datetime(2022, 1, 24), schedule_interval=None, render_template_as_native_obj=True, default_args={}, params={ "param2": "arya2@gmail.com", "sourcedir": ['/home/arya/'], "timenum": 0 }, catchup=False ) as dag: # 定义前置任务的PythonOperator python_defination_task = PythonOperator( task_id="python_defination_task", python_callable=python_defination_func, provide_context=True # 开启上下文传递,让函数能获取ti等参数 ) # 定义捕获返回值的PythonOperator value_capture_task = PythonOperator( task_id="value_capture_task", python_callable=value_capture, provide_context=True ) # 设置任务执行顺序:前置任务完成后再执行捕获任务 python_defination_task >> value_capture_task
关键修复说明
- 避免循环引用:将业务函数与Operator变量名区分开,比如用
python_defination_func作为函数名,python_defination_task作为Operator变量名,解决未定义问题。 - 上下文传递优化:Airflow 2.x支持直接将
ti作为函数参数(无需**kwargs),此时无需设置provide_context=True,示例如下:def value_capture(ti): time_info = ti.xcom_pull(task_ids='python_defination_task') print(f"前置任务返回值: {time_info}") return time_info - 任务依赖设置:必须通过
>>明确任务执行顺序,确保捕获任务在前置任务完成后执行,否则会拉取不到XCom值。 - XCom默认行为:PythonOperator会自动将函数返回值推送到XCom,无需手动调用
ti.xcom_push(),直接用ti.xcom_pull()即可拉取。
内容的提问来源于stack exchange,提问作者Arya
相关产品推荐
相关产品推荐

