如何从@task装饰器任务向Airflow传统Operator传递XCom?
问题描述
使用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中提取上游任务返回的字典。
- 修改
config任务的multiple_outputs参数为False(若需传递整个字典,无需拆分XCom):
@task(multiple_outputs=False) def config() -> dict: return { "name": "my-spark-app", }
- 修改
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

