Spark 3.5.1写入BigQuery时Decimal类型映射不符文档问题问询
解决Spark Decimal与BigQuery Numeric映射不符问题
问题根源
Dataproc 2.2.40-ubuntu22镜像默认搭载的spark-bigquery-connector版本低于0.23.2,旧版本的类型映射逻辑与新版存在差异:旧版会将精度≥30的Spark Decimal类型映射为BigQuery的BigNumeric,而0.23.2及以上版本调整了逻辑,可将Decimal(38,0)正确映射为BigQuery Numeric(符合BigQuery Numeric最大支持38位精度的特性)。
具体解决办法
替换集群默认Connector Jar
- 创建初始化脚本(例如
replace-bigquery-connector.sh):#!/bin/bash SPARK_JARS_DIR=/usr/lib/spark/jars # 删除旧版本Connector文件 rm -f ${SPARK_JARS_DIR}/spark-bigquery-*.jar # 下载指定版本的Connector gsutil cp gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.23.2.jar ${SPARK_JARS_DIR}/ - 创建集群时指定该初始化脚本:
gcloud dataproc clusters create <集群名称> \ --region <区域> \ --image-version 2.2.40-ubuntu22 \ --initialization-actions gs://<你的存储桶路径>/replace-bigquery-connector.sh - 若集群已运行,可登录主节点执行上述脚本命令,之后重启Spark服务生效。
- 创建初始化脚本(例如
提交作业时指定Connector版本
无需替换集群默认Jar,在提交Spark作业时通过参数指定新版Connector:spark-submit \ --class <你的作业主类> \ --jars gs://spark-lib/bigquery/spark-bigquery-with-dependencies_2.12-0.23.2.jar \ --master yarn \ <你的作业Jar路径>注意需确保作业依赖无版本冲突。
验证映射逻辑
可通过简单测试代码确认映射是否正常:import org.apache.spark.sql.SparkSession import org.apache.spark.sql.types.DecimalType object DecimalMappingTest { def main(args: Array[String]): Unit = { val spark = SparkSession.builder.appName("DecimalMappingTest").getOrCreate() import spark.implicits._ // 构造Decimal(38,0)类型数据 val testData = Seq((BigDecimal("12345678901234567890123456789012345678"))) val df = testData.toDF("decimal_col").select($"decimal_col".cast(DecimalType(38, 0))) // 写入BigQuery验证类型 df.write .format("bigquery") .option("table", "<项目ID>.<数据集ID>.<目标表名>") .mode("overwrite") .save() spark.stop() } }
内容的提问来源于stack exchange,提问作者Abhilash
相关产品推荐
相关产品推荐

