升级Dataproc Serverless后Spark BQ Connector分区规格不兼容问题
升级技术栈后,Spark任务写入已分区BigQuery表时触发如下错误:
Caused by: com.google.cloud.spark.bigquery.repackaged.com.google.cloud.bigquery.BigQueryException: Incompatible table partitioning specification. Expects partitioning specification none, but input partitioning specification is interval(type:day,field:event_date)
升级细节:
- 原环境:Scala 2.12.15、Spark 3.3.4、sparkBigQueryConnector 0.28.1、Dataproc Serverless 1.1、Java 11
- 新环境:Scala 2.12.15、Spark 3.5.1、sparkBigQueryConnector 0.36.4、Dataproc Serverless 1.2、Java 17
现象:删除目标表后首次运行任务可成功创建并写入数据,但再次运行(表已存在)就触发上述错误,旧版本无此问题。
新版本Spark BigQuery Connector(0.36.x及以上)对表结构校验逻辑更严格:当写入已存在的分区表时,若代码中指定了partitionField选项,连接器会尝试重新设置表的分区规则,但目标表已配置好分区,两者冲突导致报错。旧版本连接器会自动忽略该选项,直接适配已有表结构。
方案1:写入已存在的分区表时移除partitionField选项
目标表已配置分区规则,无需再通过代码指定partitionField,连接器会自动识别表的分区配置。修改写入代码如下:
dataFrame .withColumn(partitionColumn, to_date(col(partitionColumn))) .write .format(StorageOptions.BigQuery) // 修正原代码拼写错误:BiqQuery → BigQuery .mode(SaveMode.Overwrite) .option("createDisposition", "CREATE_IF_NEEDED") // 移除partitionField选项 .save(temporaryTableName)
方案2:区分表存在/不存在的场景(可选)
如果需要在表不存在时自动创建分区表,同时在表存在时适配已有结构,可以通过提前检查表是否存在分支处理:
import com.google.cloud.bigquery.BigQueryOptions import com.google.cloud.bigquery.TableId // 初始化BigQuery客户端 val bigQuery = BigQueryOptions.getDefaultInstance().getService() val tableId = TableId.of(projectId, datasetId, temporaryTableName) val tableExists = bigQuery.getTable(tableId) != null val writer = dataFrame .withColumn(partitionColumn, to_date(col(partitionColumn))) .write .format(StorageOptions.BigQuery) .mode(SaveMode.Overwrite) .option("createDisposition", "CREATE_IF_NEEDED") // 仅当表不存在时指定partitionField if (!tableExists) { writer.option("partitionField", partitionColumn) } writer.save(temporaryTableName)
若项目未引入BigQuery客户端依赖,需在build.sbt中添加:
val bigQueryClientVersion = "2.34.1" val bigQueryClient = "com.google.cloud" % "google-cloud-bigquery" % bigQueryClientVersion % "provided"
方案3:使用WRITE_TRUNCATE替代SaveMode.Overwrite(可选)
如果需要全量覆盖表数据但保留原有分区配置,可启用writeDisposition为WRITE_TRUNCATE:
dataFrame .withColumn(partitionColumn, to_date(col(partitionColumn))) .write .format(StorageOptions.BigQuery) .option("writeDisposition", "WRITE_TRUNCATE") .option("createDisposition", "CREATE_IF_NEEDED") .save(temporaryTableName)
WRITE_TRUNCATE会清空表数据但保留表结构(包括分区规则),避免分区配置冲突。
- 确保
spark.sql.sources.partitionOverwriteMode设置为DYNAMIC,配合SaveMode.Overwrite可实现动态覆盖指定分区,而非全表覆盖 - 严格遵循Dataproc Serverless官方推荐的组件版本组合,避免兼容性问题
内容的提问来源于stack exchange,提问作者Mujtaba Faizi

