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

Airflow中如何向Operator的.execute()传入context?(KeyError: 'ti')

问题解答:Airflow动态任务执行DataflowOperator报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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 06:37:44