GCP Composer2(Airflow2) Dataproc提交PySpark如何配置--packages参数
DataprocSubmitJobOperator配置Maven自动拉取依赖方案
本地spark-submit --packages的能力对应Spark原生配置项spark.jars.packages,在DataprocSubmitJobOperator中按以下方式配置即可实现完全一致的效果,无需手动上传Jar包到GCS。
正确配置示例
直接在pyspark_job的properties字段中传入对应配置,注意配置层级不要写错:
from airflow.providers.google.cloud.operators.dataproc import DataprocSubmitJobOperator PYSPARK_JOB = { "reference": {"project_id": GCP_PROJECT_ID}, "placement": {"cluster_name": DATAPROC_CLUSTER_NAME}, "pyspark_job": { "main_python_file_uri": "gs://your-code-bucket/pyspark/streaming_job.py", # 移除原有jar_file_uris配置,依赖通过Maven自动拉取 "file_uris": [ "gs://your-config-bucket/kafka/client_truststore.jks", "gs://your-config-bucket/job/prod_config.yaml" ], "properties": { # 等价于--packages参数,多个Maven坐标用英文逗号分隔,和本地提交时的传参完全一致 "spark.jars.packages": "org.apache.spark:spark-sql-kafka-0-10_2.12:3.2.0,org.mongodb.spark:mongo-spark-connector_2.12:3.0.2,com.google.cloud.spark:spark-bigquery-with-dependencies_2.12:0.27.0", # 若Dataproc集群节点无公网访问权限,必须配置可达的私有Maven仓/镜像地址,否则依赖拉取会静默失败 "spark.jars.repositories": "你的内部Maven仓库地址", # 其余Spark业务配置按需添加 "spark.sql.shuffle.partitions": "200" }, "args": ["--env", "prod"] } } submit_job = DataprocSubmitJobOperator( task_id="submit_pyspark_streaming_job", job=PYSPARK_JOB, region=REGION, project_id=GCP_PROJECT_ID )
之前配置失效的常见原因
配置后抛出java.lang.ClassNotFoundException: Failed to find data source: mongo基本是以下问题导致:
- 配置层级错误:
spark.jars.packages必须放在pyspark_job.properties路径下,若将properties字段和pyspark_job平级,配置不会被传递到Spark作业进程,完全不生效。 - 依赖拉取失败:Dataproc默认镜像不会内置第三方Maven仓地址,若集群节点没有公网访问能力,又没配置内部可达的Maven仓库,依赖拉取过程会静默失败,最终报类不存在错误。
- 版本不匹配:传入的依赖包Scala版本、Spark版本必须和Dataproc集群镜像内置的Spark、Scala版本完全对齐,否则会出现类版本不兼容,表象同样是类找不到。
生产环境注意事项
长期运行的Structured Streaming作业不建议依赖Maven自动拉取能力,Maven仓网络波动、服务不可用都会直接导致作业启动失败。生产环境建议提前将验证过的依赖Jar包上传到GCS固定路径,通过
jar_file_uris指定,稳定性更高。自动拉取依赖更适合短周期测试、临时批处理作业场景。
内容的提问来源于stack exchange,提问作者Karan Alang
相关产品推荐
相关产品推荐

