Airflow 1环境下如何向Dataproc传递参数与模板字段
TypeError: Parameter to MergeFrom() must be instance of same class: expected google.cloud.dataproc.v1beta2.OrderedJob got str 属于典型的依赖版本错配:apache-airflow-providers-google 全系列版本最低兼容要求为Airflow 2.0+、Python 3.7+,与Airflow 1.10.12内置的google-cloud-dataproc客户端、protobuf序列化库版本完全不匹配。该报错和parameters字段的使用无关,哪怕不传任何自定义参数,算子初始化构造API请求对象时就会触发类型校验失败。
Airflow 1.10.x 环境下不需要安装额外provider包,内置的contrib层Dataproc算子本身就支持向集群任务传递自定义参数(包括Jinja模板字段),具体实现方式分两类场景:
使用内置的DataprocPySparkOperator,通过算子自带的arguments字段传参,该字段原生支持Jinja模板渲染,集群上的PySpark脚本可直接通过argparse接收参数。
示例代码:
from datetime import datetime from airflow import DAG from airflow.contrib.operators.dataproc_operator import DataprocPySparkOperator default_args = { "owner": "data-team", "start_date": datetime(2024, 1, 1) } dag = DAG("dataproc_predict_dag", default_args=default_args, schedule_interval="0 2 * * *") predict_task = DataprocPySparkOperator( task_id="run_predict_job", cluster_name="predict-cluster", region="cn-north1", main_python_file_uri="gs://bucket-path/predict_code.py", # 所有自定义参数、模板字段按顺序传入即可 arguments=[ "--run_date", "{{ ds_nodash }}", "--model_version", "v3.1.0", "--input_path", "gs://bucket-path/input/{{ ds_nodash }}/" ], dag=dag )
脚本侧argparse定义和普通Python脚本一致,直接对应参数名解析即可,不需要额外适配。
如果需要按需创建临时集群、任务跑完自动删除资源,不要使用Airflow2专属的DataprocInstantiateInlineWorkflowTemplateOperator,使用内置的DataprocWorkflowTemplateInstantiateInlineOperator,在工作流任务定义的args字段传参即可,同样支持模板变量。
示例代码:
from datetime import datetime from airflow import DAG from airflow.contrib.operators.dataproc_operator import DataprocWorkflowTemplateInstantiateInlineOperator default_args = { "owner": "data-team", "start_date": datetime(2024, 1, 1) } dag = DAG("dataproc_inline_predict_dag", default_args=default_args, schedule_interval="0 2 * * *") inline_predict_task = DataprocWorkflowTemplateInstantiateInlineOperator( task_id="run_inline_predict", region="cn-north1", template={ "placement": { "managed_cluster": { "cluster_name": "temp-predict-{{ ds_nodash }}", "config": { "master_config": {"num_instances": 1, "machine_type_uri": "n1-standard-4"}, "worker_config": {"num_instances": 3, "machine_type_uri": "n1-standard-8"}, "software_config": {"image_version": "2.0-debian10"} } } }, "jobs": [ { "step_id": "predict_step", "pyspark_job": { "main_python_file_uri": "gs://bucket-path/predict_code.py", # 传参位置 "args": [ "--run_date", "{{ ds_nodash }}", "--model_version", "v3.1.0" ] } } ] }, dag=dag )
- Airflow 1.x环境禁止安装任何
apache-airflow-providers-*前缀的依赖包,该系列包为Airflow2专属,安装后会触发各类依赖冲突、序列化报错 - 上述两个算子的参数字段均已加入模板渲染列表,传入
{{ ds_nodash}}这类变量时不需要手动调用渲染函数,算子执行时会自动替换为对应值 - 传入参数时注意保持参数前缀(如
--run_date)、顺序和脚本侧argparse定义一致,避免参数解析失败
内容的提问来源于stack exchange,提问作者Imad

