从BigQuery数组列读取字段名,Spark选列遇解析错误
问题解决:Spark读取BigQuery数组列用于DataFrame列选择报错
问题根源
当前错误是因为把数组列转换成了逗号分隔的单个字符串,导致Spark将BH_TCHH,BH_TCHF当作一个完整的列名去查找,而非识别成两个独立列。具体问题出在两处:
- 读取BigQuery时用
concat_ws(",", $"column_list")把数组转成字符串,丢失了原生数组结构; - 后续用
columnList.mkString(",")再次将序列转成单个字符串,传递给select时被当作单一列名。
解决方案
1. 正确读取BigQuery的数组列
保留数组原生结构,直接读取为Seq[String],不要做字符串拼接:
// 读取BigQuery表,直接获取column_list数组列 val columnList: Seq[String] = spark.read .format("bigquery") .option("table", s"$confProjectId.$mappingsDatasetName.$busyHourTableName") .option("project", confProjectId) .load() .select($"column_list") .where(col("TECHNOLOGY") === technology && col("CGNAME") === cgName && col("REGIONAL_UNIT") === regionalUnit) .as[Seq[String]] // 直接将数组列映射为Seq[String] .collect() .head // 假设where条件仅返回一行,取第一个匹配结果
2. 正确使用数组列名进行DataFrame选择
将Seq[String]类型的列名列表转换为Spark的Column对象序列,与其他列合并后传入select:
// 将各列名转换为Column对象,合并后传入select val selectedColumns = Seq(col(timestampColumnName)) ++ identifierColumnNames.map(col) ++ columnList.map(col) val filtereddfRaw = dfRaw.select(selectedColumns: _*)
关键说明
- Spark的
select方法接受多个Column对象或列名字符串作为参数,但如果传入逗号分隔的字符串,它会将其视为单个列名而非多个列; - 通过
map(col)将字符串列名转换为Column对象,再用++合并多个列名序列,最后用: _*将序列展开为可变参数传递给select,即可正确选择目标列。
内容的提问来源于stack exchange,提问作者Vikrant Singh Rana
相关产品推荐
相关产品推荐

