Databricks Runtime15.4能否替换内置BigQuery连接器为0.43.x版本?
我在AWS上使用Databricks Runtime 15.4(Spark 3.5 / Scala 2.12),目标是使用最新的Google BigQuery连接器的direct写入方法(基于BigQuery Storage Write API)——该方法无需临时GCS桶,是我的环境必需的:
option("writeMethod", "direct")
我通过Maven安装了官方Google连接器作为集群库:
com.google.cloud.spark:spark-bigquery-with-dependencies_2.12:0.43.1
库安装成功并显示为“已附加”,但Databricks运行时并未使用该连接器。我通过以下代码检查实际加载的连接器:
jvm = spark._jvm provider = jvm.com.google.cloud.spark.bigquery.BigQueryRelationProvider() location = provider.getClass().getProtectionDomain().getCodeSource().getLocation().toString() print(location)
输出始终为:
...spark-bigquery-connector-hive-2.3__hadoop-3.2_2.12--fatJar-assembly-0.22.2-SNAPSHOT.jar
这说明Databricks始终加载内置的分叉连接器(0.22.2-SNAPSHOT),而非我安装的Google连接器(0.43.x)。
其他观察结果:
- 重启集群无任何变化;
- 已安装的连接器显示为“已附加”,但从未出现在
/databricks/jars目录中; /databricks/jars仅包含:spark-bigquery-connector-hive-2.3__hadoop-3.2_2.12--fatJar-assembly-0.22.2-SNAPSHOT.jarspark-bigquery-with-dependencies_2.12-0.41.0.jar(Databricks自有副本)
spark.read.format("bigquery")每次仍解析为内置连接器。
核心问题:在Databricks Runtime 15.4上是否有受支持的方法覆盖或替换内置BigQuery连接器,使spark.read.format("bigquery")使用Google官方的spark-bigquery-with-dependencies_2.12(0.43.x)连接器,以实现无需临时GCS桶的direct write方法?或者Databricks的BigQuery连接器版本是固定且不可由用户覆盖的?
方法1:通过Spark配置强制指定连接器类
Databricks内置连接器会优先注册bigquery格式,你可以在集群的Spark配置中添加以下参数,强制Spark使用Google官方连接器的实现类:
spark.sql.extensions=com.google.cloud.spark.bigquery.BigQuerySparkExtensions spark.sql.sources.bigquery=com.google.cloud.spark.bigquery.BigQueryRelationProvider
添加后重启集群,此时spark.read.format("bigquery")会加载你安装的官方连接器。
方法2:不依赖bigquery格式,直接指定连接器类
如果不想修改全局配置,可以在读写时直接指定官方连接器的类名,绕过内置的格式映射:
# 读取数据 df = spark.read.format("com.google.cloud.spark.bigquery")\ .option("project", "your-project-id")\ .option("dataset", "your-dataset")\ .option("table", "your-table")\ .load() # 写入数据(使用direct方法) df.write.format("com.google.cloud.spark.bigquery")\ .option("writeMethod", "direct")\ .option("project", "your-project-id")\ .option("dataset", "your-dataset")\ .option("table", "your-table")\ .mode("append")\ .save()
方法3:使用集群初始化脚本替换内置Jar
如果上述方法无效,可以通过集群初始化脚本将官方连接器Jar复制到/databricks/jars目录,覆盖内置的Jar文件:
- 创建一个初始化脚本(例如
replace_bq_connector.sh):
#!/bin/bash # 找到安装的官方连接器Jar路径(通常在/dbfs/FileStore/jars下) OFFICIAL_JAR=$(find /dbfs/FileStore/jars -name "spark-bigquery-with-dependencies_2.12-0.43.1.jar" -type f) # 复制到/databricks/jars目录,替换内置的Jar cp $OFFICIAL_JAR /databricks/jars/ # 删除内置的分叉连接器Jar rm /databricks/jars/spark-bigquery-connector-hive-2.3__hadoop-3.2_2.12--fatJar-assembly-0.22.2-SNAPSHOT.jar
- 将脚本上传到DBFS(例如
dbfs:/scripts/replace_bq_connector.sh) - 在集群配置中添加初始化脚本路径,重启集群。
注意事项
- 确保你安装的官方连接器版本与Databricks Runtime的Spark、Scala版本兼容(0.43.1适配Spark 3.5/Scala 2.12,符合你的环境);
- 使用初始化脚本替换Jar可能会影响Databricks内置的BigQuery集成功能,建议先在测试集群验证;
- Databricks官方并未完全支持覆盖内置连接器,因此上述方法属于变通方案,后续Runtime版本更新可能需要调整。
内容的提问来源于stack exchange,提问作者Thilina

