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

如何在Airflow中复用PythonOperator并传入不同XCom参数?

解决方法

不需要传入整个task_instance作为参数,更简洁的方式是给function_two增加一个参数,指定要拉取XCom的目标任务ID,这样就能复用同一个函数处理不同任务返回的列表。

步骤1:修改function_two,支持动态指定任务ID

把硬编码的任务ID改成可传入的参数,同时保留获取task_instance的逻辑:

def function_two(target_task_id, **kwargs):
    # 从上下文获取task_instance(Airflow中常用ti作为缩写)
    ti = kwargs['ti']
    # 根据传入的任务ID拉取对应的XCom值
    my_list = ti.xcom_pull(task_ids=target_task_id)
    # 这里写你的列表处理逻辑
    print(f"处理来自任务{target_task_id}的列表:{my_list}")

步骤2:创建生成列表的任务

比如你的function_one和new_function对应的Operator:

def function_one(**kwargs):
    return [1, 2, 3]

def new_function(**kwargs):
    return [4, 5, 6]

# 生成第一个列表的任务
get_one = PythonOperator(
    task_id='one',
    python_callable=function_one,
    provide_context=True
)

# 生成第二个列表的任务
get_new = PythonOperator(
    task_id='new_task',
    python_callable=new_function,
    provide_context=True
)

步骤3:复用function_two处理不同列表

创建两个PythonOperator,分别传入不同的目标任务ID:

# 处理来自task_id='one'的列表
process_list_one = PythonOperator(
    task_id='process_one',
    python_callable=function_two,
    op_args=['one'],  # 传入要拉取的任务ID
    provide_context=True  # 必须开启,才能让function_two获取到task_instance
)

# 处理来自task_id='new_task'的列表
process_list_new = PythonOperator(
    task_id='process_new',
    python_callable=function_two,
    op_args=['new_task'],  # 传入另一个任务ID
    provide_context=True
)

步骤4:设置任务依赖

最后把任务按顺序串联起来:

get_one >> process_list_one
get_new >> process_list_new

这样就实现了用同一个function_two处理不同任务返回的列表,完全不需要修改函数内部的逻辑,只需要通过参数传递目标任务ID即可。

内容的提问来源于stack exchange,提问作者x89

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 17:10:33