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
相关产品推荐
相关产品推荐

