Airflow中如何向Operator的.execute()传入context?(KeyError: 'ti')
KeyError: 'ti' 背景
Airflow版本=2.3.4,Python版本=3.8,在Google Cloud Composer中开发DAG,需求是根据Google Cloud Storage桶中的文件数量,实例化对应数量的DataflowOperator任务。
错误现象与日志
执行DAG时,虽能根据列表长度创建对应数量的映射实例、启动对应数量的Dataflow作业,但最终任务及整个DAG失败,抛出以下错误:
... File "/opt/python3.8/lib/python3.8/site-packages/airflow/models/baseoperator.py", line 1390, in xcom_push context['ti'].xcom_push(key=key, value=value, execution_date=execution_date) KeyError: 'ti' {taskinstance.py:1408} INFO - Marking task as FAILED. ...
问题代码示例
... with models.DAG( DAG_NAME, default_args=DEFAULT_ARGS, schedule_interval=None, start_date=datetime.now(), catchup=False, ) as dag: start = dummy.DummyOperator(task_id='start', trigger_rule='all_success') @task def dynamic_task_test(sample_list): dataflow_operator = DataflowStartFlexTemplateOperator( task_id = 'dataflow_task_id', project_id = 'project_id', location = 'location', body = { 'some_parameters' : f'{sample_list}'}, ) dataflow_operator.execute(dict()) start >> dynamic_task_test.expand(sample_list=['A', 'B', 'C'])
问题解答
1. 错误代码及信息的含义
KeyError: 'ti' 表示Airflow的Operator在执行xcom_push操作时,找不到上下文(context)中的ti(TaskInstance,任务实例)对象。原因是你直接调用了Operator的execute()方法,但传入的是空字典,没有传递Airflow任务执行所需的核心上下文参数;而DataflowStartFlexTemplateOperator内部会尝试用XCom传递数据,依赖上下文里的ti对象,因此触发报错。
2. 正确解决方法
不要在@task装饰的函数内部实例化并调用Operator的execute()方法,直接利用Airflow 2.x TaskFlow API的动态任务映射特性,对DataflowStartFlexTemplateOperator使用partial()+expand()组合来生成动态任务:
... with models.DAG( DAG_NAME, default_args=DEFAULT_ARGS, schedule_interval=None, start_date=datetime.now(), catchup=False, ) as dag: start = dummy.DummyOperator(task_id='start', trigger_rule='all_success') # 用partial固定公共参数,expand动态生成不同body的任务实例 dataflow_task = DataflowStartFlexTemplateOperator.partial( task_id='dataflow_task_id', project_id='project_id', location='location', ).expand( body=[{'some_parameters': 'A'}, {'some_parameters': 'B'}, {'some_parameters': 'C'}] ) start >> dataflow_task
这种方式下Airflow会自动处理上下文传递,避免手动调用execute()带来的问题。
3. 向Operator的.execute()方法正确传入context
如果确实需要手动调用execute()(不推荐这种用法),需在@task装饰的函数中接收context参数,并将其传递给execute():
@task def dynamic_task_test(sample_list, context): dataflow_operator = DataflowStartFlexTemplateOperator( task_id='dataflow_task_id', project_id='project_id', location='location', body={'some_parameters': f'{sample_list}'}, ) dataflow_operator.execute(context)
但这种方式不符合Airflow TaskFlow API的设计规范,容易引发其他上下文相关问题,不建议使用。
4. 其他实现思路建议
- 先通过Sensor获取GCS文件列表:如果需要根据GCS桶的实际文件数量生成任务,可以先用自定义函数或
GCSSensor变种获取桶内文件列表,再传递给动态任务映射:@task def get_gcs_files(): from google.cloud import storage client = storage.Client() bucket = client.get_bucket('your-bucket-name') blobs = bucket.list_blobs() return [blob.name for blob in blobs] file_list = get_gcs_files() dataflow_task = DataflowStartFlexTemplateOperator.partial( task_id='dataflow_task', project_id='project_id', location='location' ).expand( body=[{'some_parameters': file} for file in file_list] ) - 避免使用SubDAG:SubDAG维护成本高,TaskFlow API的动态任务映射是更简洁高效的方案。
- 封装自定义Operator:如果有特殊业务逻辑,可以封装自定义Operator,但对于简单的动态任务场景,官方的动态映射已经足够满足需求。
内容的提问来源于stack exchange,提问作者pfhermosa

