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

Airflow:Python函数与TaskGroup间如何通过XCom传递变量?

问题:Airflow中Python函数与TaskGroup间的变量传递问题

我是Airflow新手,需要在Python函数和TaskGroup之间交换变量(知道这并非Airflow的核心用途,但场景必须这么做)。以下是我的代码片段:

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.decorators import task_group
import logging
from pendulum import datetime

def push_function(**kwargs):
    files = ['a','b','c']
    kwargs['ti'].xcom_push(key='files', value=files)

with DAG(
    "tst",
    start_date=datetime(2023, 11, 7),
    schedule_interval="30 6 * * *",
    catchup=False,
) as dag:
    push = PythonOperator(
        task_id='push_task',
        python_callable=push_function,
        dag=dag,
    )

    @task_group(group_id="group")
    def pull_task(**kwargs):
        data= kwargs['ti'].xcom_pull(task_ids='push_task', key='files')
        logging.info(f"transfaired Variable: {folders_today}")
        for item in data:
                filepath = f"/tmp/{item}.xml"
                extract_load = SFTPOperator(
                    task_id=f"download_{item}",
                    ssh_conn_id="sftp",
                    remote_filepath=f"{item}.xml",
                    local_filepath=filepath,
                    operation="get",
                    create_intermediate_dirs=True
)
      
    push >> pull_task()

核心问题:能否将push_function中的变量files传递到pull_task中?

当前执行语句:

folders_today= kwargs['ti'].xcom_pull(task_ids='push_task', key='files')

时出现错误:

KeyError: 'ti'

了解到TaskGroup的上下文似乎无法在该函数中使用,但不清楚具体含义和如何让上下文可用。


原因分析

@task_group装饰的函数是用于定义TaskGroup结构的构建函数,它在DAG解析阶段运行(即Airflow加载DAG文件时),而非任务执行阶段。这个阶段不存在任务实例(ti),自然无法获取到运行时上下文,所以直接在这里调用xcom_pull会触发KeyError。

解决方案

要实现需求,必须把XCom拉取的逻辑放到任务执行阶段运行的PythonOperator中,再通过这个任务动态生成SFTPOperator。以下是两种可行的实现方式:

方式1:用PythonOperator作为TaskGroup的入口任务

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.decorators import task_group
from airflow.providers.sftp.operators.sftp import SFTPOperator
import logging
from pendulum import datetime

def push_function(**kwargs):
    files = ['a','b','c']
    kwargs['ti'].xcom_push(key='files', value=files)

def generate_download_tasks(**kwargs):
    # 这里能获取到ti,因为是PythonOperator的执行函数,运行在任务执行阶段
    data = kwargs['ti'].xcom_pull(task_ids='push_task', key='files')
    logging.info(f"transferred Variable: {data}")
    tasks = []
    for item in data:
        filepath = f"/tmp/{item}.xml"
        extract_load = SFTPOperator(
            task_id=f"download_{item}",
            ssh_conn_id="sftp",
            remote_filepath=f"{item}.xml",
            local_filepath=filepath,
            operation="get",
            create_intermediate_dirs=True
        )
        tasks.append(extract_load)
    return tasks

with DAG(
    "tst",
    start_date=datetime(2023, 11, 7),
    schedule_interval="30 6 * * *",
    catchup=False,
) as dag:
    push = PythonOperator(
        task_id='push_task',
        python_callable=push_function,
        dag=dag,
    )

    @task_group(group_id="group")
    def pull_task_group():
        # 在TaskGroup内定义PythonOperator,负责拉取XCom并生成子任务
        generate_tasks = PythonOperator(
            task_id="generate_download_tasks",
            python_callable=generate_download_tasks,
            provide_context=True  # Airflow 2.x默认开启,显式声明更清晰
        )
        return generate_tasks

    push >> pull_task_group()

方式2:用TaskGroup上下文管理器实现

from airflow import DAG
from airflow.operators.python import PythonOperator
from airflow.utils.task_group import TaskGroup
from airflow.providers.sftp.operators.sftp import SFTPOperator
import logging
from pendulum import datetime

def push_function(**kwargs):
    files = ['a','b','c']
    kwargs['ti'].xcom_push(key='files', value=files)

with DAG(
    "tst",
    start_date=datetime(2023, 11, 7),
    schedule_interval="30 6 * * *",
    catchup=False,
) as dag:
    push = PythonOperator(
        task_id='push_task',
        python_callable=push_function,
        dag=dag,
    )

    # 使用TaskGroup上下文管理器替代装饰器
    with TaskGroup(group_id="group") as pull_task_group:
        def pull_and_generate(**kwargs):
            data = kwargs['ti'].xcom_pull(task_ids='push_task', key='files')
            logging.info(f"transferred Variable: {data}")
            tasks = []
            for item in data:
                filepath = f"/tmp/{item}.xml"
                task = SFTPOperator(
                    task_id=f"download_{item}",
                    ssh_conn_id="sftp",
                    remote_filepath=f"{item}.xml",
                    local_filepath=filepath,
                    operation="get",
                    create_intermediate_dirs=True
                )
                tasks.append(task)
            return tasks
        
        generate_tasks = PythonOperator(
            task_id="generate_download_tasks",
            python_callable=pull_and_generate,
            provide_context=True
        )

    push >> pull_task_group

关键要点

  • TaskGroup的构建函数(装饰器或上下文管理器内的顶层代码)仅负责定义任务结构,不能处理运行时数据(如XCom)。
  • 所有需要访问运行时上下文(ti、XCom等)的逻辑,必须放在PythonOperator或其他执行类任务的函数中。
  • 动态生成的子任务会自动归属到所在的TaskGroup,无需额外配置。

内容的提问来源于stack exchange,提问作者Moriz Bühler

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 10:45:11