Airflow中如何向DataprocSubmitJobOperator传递参数?
问题解决:Airflow DataprocSubmitJobOperator向Spark作业传参的正确方式
你之前的配置错误点在于将业务自定义参数直接作为spark_job的顶层字段填写,Dataproc API不会识别这些自定义字段,也不会将其传递给Spark主类的入参列表。
你用gcloud命令提交时,写在命令末尾的gcsPath=xxx这类键值对,本质是传递给主类main方法的参数,对应到Dataproc的作业配置里,需要放在spark_job下的args数组字段中。另外你gcloud命令里配置的Spark运行参数(--properties后的内容),对应要放在spark_job下的properties字段里。
正确配置示例
# SPARK JOB CONF SPARK_JOB = { "reference": {"project_id": PROJECT_ID}, "placement": {"cluster_name": CLUSTER_NAME}, "spark_job": { "jar_file_uris": ["gs://testing-bucket/test-1.1.0.jar"], "main_class": "com.walmart.ei.testClass", # 对应gcloud命令的--properties配置项 "properties": { "spark.ui.killEnabled": "true", "spark.dynamicAllocation.enabled": "true", "spark.ui.enabled": "true", "spark.master": "yarn" }, # 传递给Spark主类的自定义参数,每个键值对为数组的一个元素 "args": [ "gcsPath=gs://testing_part/path/", "test.id=404", "correlation.id=AUDIT_2222", "b.url=http://testing-v1.net/test/" ] } } spark_task = DataprocSubmitJobOperator( task_id="spark_task", job=SPARK_JOB, location=REGION, project_id=PROJECT_ID, gcp_conn_id=CONN_ID )
验证逻辑
提交作业后可以查看Spark作业的stdout日志,确认代码中println(props)的输出是否包含所有传入的键值对,即可验证参数是否传递成功。
内容的提问来源于stack exchange,提问作者abhimanyu
相关产品推荐
相关产品推荐

