Airflow DataprocSubmitJobOperator报错:Job无python_file_uris字段
Airflow DataprocSubmitJobOperator传递PyFiles报错修复
错误原因
你将python_file_uris字段配置在了PYSPARK_JOB的顶层字典中,但根据Dataproc的Job协议定义,这个字段属于pyspark_job子字典的专属配置项,顶层Job结构体并不包含该字段,因此触发了ValueError: Protocol message Job has no "python_file_uris" field报错。
修复步骤
修改PYSPARK_JOB的结构,将python_file_uris移入pyspark_job内部即可解决问题,修改后的代码如下:
PYSPARK_JOB = { "reference": {"project_id": PROJECT_ID}, "placement": {"cluster_name": CLUSTER_NAME}, "pyspark_job": { "main_python_file_uri": PYSPARK_URI, "jar_file_uris" : ["gs://dataproc-spark-jars/mongo-spark-connector_2.12-3.0.2.jar", 'gs://dataproc-spark-jars/bson-4.0.5.jar','gs://dataproc-spark-jars/mongo-spark-connector_2.12-3.0.2.jar','gs://dataproc-spark-jars/mongodb-driver-core-4.0.5.jar', 'gs://dataproc-spark-jars/mongodb-driver-sync-4.0.5.jar','gs://dataproc-spark-jars/spark-avro_2.12-3.1.2.jar','gs://dataproc-spark-jars/spark-bigquery-with-dependencies_2.12-0.23.2.jar', 'gs://dataproc-spark-jars/spark-token-provider-kafka-0-10_2.12-3.1.3.jar','gs://dataproc-spark-jars/htrace-core4-4.1.0-incubating.jar','gs://dataproc-spark-jars/hadoop-client-3.3.1.jar','gs://dataproc-spark-jars/spark-sql-kafka-0-10_2.12-3.1.3.jar','gs://dataproc-spark-jars/hadoop-client-runtime-3.3.1.jar','gs://dataproc-spark-jars/hadoop-client-3.3.1.jar','gs://dataproc-spark-jars/kafka-clients-3.2.0.jar','gs://dataproc-spark-jars/commons-pool2-2.11.1.jar'], "file_uris":['gs://kafka-certs/versa-kafka-gke-ca.p12','gs://kafka-certs/syslog-vani.p12', 'gs://kafka-certs/alarm-compression-user.p12','gs://kafka-certs/appstats-user.p12', 'gs://kafka-certs/insights-user.p12','gs://kafka-certs/intfutil-user.p12', 'gs://kafka-certs/reloadpred-chkpoint-user.p12','gs://kafka-certs/reloadpred-user.p12', 'gs://dataproc-spark-configs/topic-customer-map.cfg','gs://dataproc-spark-configs/params.cfg','gs://kafka-certs/issues-user.p12','gs://kafka-certs/anomaly-user.p12','gs://kafka-certs/appstat-anomaly-user.p12','gs://kafka-certs/appstat-agg-user.p12','gs://kafka-certs/alarmblock-user.p12'], # 将python_file_uris移至pyspark_job内部 "python_file_uris": ['gs://dagger-mongo/move2mongo_api.zip'] } }
额外优化建议
- 配置
python_file_uris后,Dataproc会自动将指定的zip包分发到集群所有节点的Spark运行环境中,无需在任务代码中手动调用spark.sparkContext.addPyFile("move2mongo_api.zip"),可以移除该行代码避免重复操作。 - 确认
move2mongo_api.zip的GCS路径正确,且Dataproc集群拥有访问该存储桶的权限(可通过IAM角色配置)。
内容的提问来源于stack exchange,提问作者Karan Alang
相关产品推荐
相关产品推荐

