在Google Dataproc集群通过Airflow运行Spark Pipeline时遇类未找到异常
问题根源
org.apache.spark.sql.connector.read.SupportsRuntimeFiltering是Spark 3.3及以上版本才引入的API类。错误日志显示,你使用的BigQuery连接器(spark-3.3-bigquery-0.30.0.jar)在调用这个类,但当前Dataproc集群的Spark版本可能低于3.3;另外你手动添加的高版本Spark核心jar(spark-sql_2.13-3.4.0.jar)与集群自带的Spark版本冲突,导致类加载失败。
Dataproc集群已预安装对应版本的完整Spark环境,手动上传不同版本的核心jar会干扰原有依赖,必然引发类加载异常。
解决方案
1. 确认集群Spark版本
登录集群主节点,执行命令查看当前Spark版本:
spark-submit --version
或查看集群安装目录验证:
ls /usr/lib/spark/
2. 匹配BigQuery连接器与Spark版本
BigQuery连接器版本必须与集群Spark版本严格对应:
- 集群为Spark 3.2:使用
spark-3.2-bigquery-xxx.jar - 集群为Spark 3.3+:使用
spark-3.3-bigquery-xxx.jar(无需手动添加spark-sql jar,集群自带)
推荐方式:不要手动上传jar,通过spark.jars.packages参数让Spark自动拉取匹配版本的连接器(包含所有依赖),修改Dataproc Submit Job Operator配置:
pyspark_job={ "reference": {"project_id": PROJECT_ID}, "placement": {"cluster_name": CLUSTER_NAME}, "pyspark_job": { "main_python_file_uri": "gs://bucket-name/dependencies/main.py", "python_file_uris":[f"gs://bucket-name/dependencies/config.py", f"gs://bucket-name/dependencies/util.py"], "properties": { "spark.jars.packages": "com.google.cloud.spark:spark-bigquery-with-dependencies_2.12:0.30.0" } } }
注意:_2.12为Scala版本,需与集群Scala版本匹配(Dataproc Spark 3.x默认使用Scala 2.12)。
3. 移除手动添加的Spark核心jar
删除初始化脚本中下载spark-sql_2.13-3.4.0.jar的步骤,同时移除PYSPARK_SUBMIT_ARGS和Dataproc Job配置中的该jar引用。
4. 统一Scala版本兼容性
你的配置中混用了Scala 2.13(spark-sql_2.13-3.4.0.jar、scala-library-2.13.1.jar)和Scala 2.12(s3mock_2.12-0.1.8.jar)的jar,会引发类加载问题。确保所有依赖的Scala版本与集群一致。
内容的提问来源于stack exchange,提问作者Rohan Aswani

