PySpark写入BigQuery遇数据源找不到异常的问题求助
问题场景
我尝试将PySpark计算结果写入BigQuery,按指引初始化SparkSession:
from pyspark.sql import SparkSession spark = SparkSession.builder\ .config("spark.jars.packages",\ "com.google.cloud.spark:spark-bigquery-with-dependencies_2.12:0.27.1")\ .getOrCreate() spark_context = spark.sparkContext
执行保存操作时:
data.toDF(schema) \ .write.format("bigquery") \ .option("table", "tmp-project:tmpdataset.tmp_table") \ .save()
每次都会抛出异常:
java.lang.ClassNotFoundException: Failed to find data source: bigquery. Please find packages at http://spark.apache.org/third-party-projects.html
已尝试的无效方案
- 直接配置GCS路径
gs://spark-lib/bigquery/spark-bigquery-latest_2.12.jar - 本地下载
spark-bigquery-latest_2.12.jar并配置本地路径(日志显示文件存在) - 暂无法使用
pyspark --jars gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.27.1.jar方式传入依赖 - 将format从"bigquery"改为"com.google.cloud.spark.bigquery"
环境信息
- PySpark版本:3.0.0
- Scala版本:2.12.10
执行以下代码返回空结果:
[spark_context._jsc.sc().jars().apply(i) for i in range(jvc.sc().jars().length())]
更新情况
将Spark升级至3.1.1,使用本地下载并赋予chmod 777权限的jar包后,情况有变化但问题仍未解决:
- 执行
spark_context._jsc.sc().listJars()返回Vector(spark://<...>.svc:<port>/jars/spark-bigquery-with-dependencies_2.13-0.27.1.jar) - 执行
spark_context._jsc.sc().jars()返回ArrayBuffer(./<...>/spark-bigquery-with-dependencies_2.13-0.27.1.jar) - 日志显示:
SparkContext: Added JAR ./<...>/spark-bigquery-with-dependencies_2.13-0.27.1.jar at spark://<...>.svc:<port>/jars/spark-bigquery-with-dependencies_2.13-0.27.1.jar with timestamp <timestamp>
解决思路与方案
1. 严格匹配Scala版本与jar包版本
你升级后用的是spark-bigquery-with-dependencies_2.13-0.27.1.jar,但你的Scala版本是2.12.10,Scala版本必须和jar包后缀的版本号严格对应(2.12对应_2.12,2.13对应_2.13)。立即更换为对应Scala 2.12的jar包:spark-bigquery-with-dependencies_2.12-0.27.1.jar。
2. 正确配置本地jar包路径
如果使用本地jar包,不要用spark.jars.packages,改用spark.jars指定绝对路径,避免相对路径加载失败:
spark = SparkSession.builder\ .config("spark.jars", "/绝对路径/spark-bigquery-with-dependencies_2.12-0.27.1.jar")\ .getOrCreate()
3. 使用完整数据源类名
保存时的format必须写完整的数据源类名,不要用简写:
data.toDF(schema) \ .write.format("com.google.cloud.spark.bigquery") \ .option("table", "tmp-project:tmpdataset.tmp_table") \ .option("parentProject", "你的GCP项目ID") \ .save()
添加parentProject参数可以明确指定GCP项目,避免权限或路由模糊问题。
4. 验证jar包是否真正加载
执行以下代码确认Spark已加载目标jar包:
print([spark_context._jsc.sc().jars().apply(i) for i in range(spark_context._jsc.sc().jars().length())])
如果输出中能看到正确的2.12版本jar包路径,说明加载成功;如果看不到,检查路径是否正确、jar包权限是否为可读取状态(chmod 644即可,无需777)。
5. 集群环境额外检查
如果是集群运行,确保所有节点都能访问到该jar包——要么放在共享存储,要么每个节点本地都有相同路径的jar包。
内容的提问来源于stack exchange,提问作者VVildVVolf

