升级Dataproc/Spark/Scala后读取BigQuery反序列化类型不匹配问题
解决Spark 3.5 + Dataproc 2.2读取BigQuery时的反序列化类型不匹配问题
问题根源
Spark 3.5配套的BigQuery连接器新增了默认行为:将BigQuery中ARRAY<STRUCT<key STRING, value STRING>>类型的字段自动解析为Map<String, String>,而你的样例类中对应字段定义为Seq[DataUnitFunction],导致反序列化时出现类型不匹配错误:
Exception in thread "main" org.apache.spark.sql.AnalysisException: [UNSUPPORTED_DESERIALIZER.DATA_TYPE_MISMATCH] The deserializer is not supported: need a(n) "ARRAY" field but got "MAP<STRING, STRING>".
解决方案
只需在读取BigQuery表时添加parseKeyValueStructsAsMaps=false的选项,禁用自动转换key-value结构体数组为Map的行为,即可恢复与原有样例类的兼容性。
修改读取代码(局部配置)
修改getDatasetConfigurations方法中的读取逻辑,添加该选项:
def getDatasetConfigurations( spark: SparkSession, confProjectId: String, mappingsDatasetName: String, datasetConfigurationsTableName: String, ): Seq[DatasetConfigurationRow] = { import org.apache.spark.sql.functions._ import spark.implicits._ spark.read .format("bigquery") .option("parseKeyValueStructsAsMaps", "false") // 新增此选项 .option("table", s"$confProjectId.$mappingsDatasetName.$datasetConfigurationsTableName") .option("project", confProjectId) .load() .select( col("technology"), col("name"), col("identifier_column_names"), col("column_mappings"), col("timestamp_column_name")) .as[DatasetConfigurationRow] .collect() }
或全局配置(所有BigQuery读取生效)
如果需要所有BigQuery读取都保持原有行为,可以在构建SparkSession时添加全局配置:
implicit val spark: SparkSession = SparkSession.builder.appName(jobName) .config("spark.sql.bigquery.parseKeyValueStructsAsMaps", false) // 全局禁用自动转换 .master(if (runLocally) "local[2]" else "") .getOrCreate()
说明
该配置项是Dataproc 2.2/Spark 3.5对应的BigQuery连接器新增的特性,默认值为true。关闭后,连接器会将BigQuery中的结构体数组解析为Spark数组类型,与你原有样例类中的Seq[DataUnitFunction]类型完全匹配,无需修改样例类或大规模重构代码。
内容的提问来源于stack exchange,提问作者Vikrant Singh Rana
相关产品推荐
相关产品推荐

