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

如何从@task装饰器任务向Airflow传统Operator传递XCom?

Airflow TaskFlow API传递XCom到EmrServerlessCreateApplicationOperator失败的解决方法

问题描述

使用Airflow TaskFlow API的@task装饰器生成配置字典,传递给EmrServerlessCreateApplicationOperator时执行失败,直接硬编码字典则可成功运行。

报错信息

ERROR - Failed to execute job 730 for task create_spark_app_fails (botocore.client.ClientCreator._create_api_method.._api_call() argument after ** must be a mapping, not PlainXComArg; 2990)

复现代码

import datetime
from airflow.decorators import task, dag

from airflow.providers.amazon.aws.operators.emr import EmrServerlessCreateApplicationOperator
    
@dag(
    dag_id="demo-xcom-problem",
    start_date=datetime.datetime(2021, 1, 1),
    catchup=False
)
def taskflow():
    @task(multiple_outputs=True)
    def config() -> dict:
        return {
            "name": "my-spark-app",
        }
    create_app_fails = EmrServerlessCreateApplicationOperator(
            task_id="create_spark_app_fails",
            job_type="SPARK",
            release_label="emr-6.9.0",
            config=config(),
            aws_conn_id="",
        )
    create_app_succeeds = EmrServerlessCreateApplicationOperator(
            task_id="create_spark_app_succeeds",
            job_type="SPARK",
            release_label="emr-6.9.0",
            config={
                "name": "my-spark-app",
            },
            aws_conn_id="",
        )
    create_app_fails
    create_app_succeeds
taskflow()

问题原因

@task装饰的任务返回的是PlainXComArg对象,这是Airflow标记XCom传递的占位符。而EmrServerlessCreateApplicationOperator属于传统Airflow Operator,执行阶段会直接将该占位符对象传递给AWS boto3 API,但boto3要求参数为实际的字典(mapping类型),因此触发类型错误。

解决方案

方案1:使用Jinja模板渲染XCom值

利用Airflow的模板渲染功能,直接从XCom中提取上游任务返回的字典。

  1. 修改config任务的multiple_outputs参数为False(若需传递整个字典,无需拆分XCom):
@task(multiple_outputs=False)
def config() -> dict:
    return {
        "name": "my-spark-app",
    }
  1. 修改create_app_fails的config参数为Jinja模板表达式:
create_app_fails = EmrServerlessCreateApplicationOperator(
        task_id="create_spark_app_fails",
        job_type="SPARK",
        release_label="emr-6.9.0",
        config="{{ ti.xcom_pull(task_ids='config') }}",
        aws_conn_id="",
    )

Airflow会在任务执行时自动渲染该模板,将XCom中的字典值传递给config参数。

方案2:用TaskFlow包装传统Operator

将EmrServerlessCreateApplicationOperator包装在@task装饰的函数中,利用TaskFlow自动解析XComArg的能力。

修改后的完整代码:

import datetime
from airflow.decorators import task, dag

from airflow.providers.amazon.aws.operators.emr import EmrServerlessCreateApplicationOperator
    
@dag(
    dag_id="demo-xcom-fixed",
    start_date=datetime.datetime(2021, 1, 1),
    catchup=False
)
def taskflow():
    @task(multiple_outputs=False)
    def config() -> dict:
        return {
            "name": "my-spark-app",
        }

    @task
    def create_spark_app(config_dict):
        # 包装传统Operator并执行
        operator = EmrServerlessCreateApplicationOperator(
            task_id="create_spark_app_fixed",
            job_type="SPARK",
            release_label="emr-6.9.0",
            config=config_dict,
            aws_conn_id="",
        )
        return operator.execute(context=None)

    # 定义任务依赖
    config_task = config()
    create_spark_app(config_task)

taskflow()

这种方式下,TaskFlow会自动将上游的PlainXComArg解析为实际的字典,传递给包装后的任务,再传给Operator。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.26 04:44:58