You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

在Google Dataproc集群通过Airflow运行Spark Pipeline时遇类未找到异常

解决Dataproc中Spark Pipeline的ClassNotFoundException(org.apache.spark.sql.connector.read.SupportsRuntimeFiltering)

问题根源

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.19 21:23:14