Spark Scala:如何将含多记录的嵌套列拆分为多列
解决方案:Spark Scala将数组嵌套结构体转为多列
步骤1:将数组转为Map类型
先把risk_table数组转换成以index为键、整个结构体为值的Map,方便后续按索引提取字段:
import org.apache.spark.sql.functions.{col, map_from_entries, transform} val dfWithMap = df.withColumn( "risk_table_map", map_from_entries( transform(col("risk_table"), struct => struct(struct("index"), struct)) ) )
步骤2:提取目标列(固定索引场景)
如果已知需要提取的index名称(比如mean、count、min等),直接通过Map键值对提取对应字段并命名:
val finalDf = dfWithMap .select( col("symbol_id"), col("risk_table_map")("mean")("risk_buy").alias("mean"), col("risk_table_map")("count")("risk_buy").alias("count"), col("risk_table_map")("min")("risk_buy").alias("min"), col("risk_table_map")("max")("risk_buy").alias("max"), col("risk_table_map")("50%")("risk_buy").alias("50%"), col("risk_table_map")("75%")("risk_buy").alias("75%"), // 根据实际需求提取reward_buy/reward_sell,示例以count对应的字段为例 col("risk_table_map")("count")("reward_buy").alias("reward_buy"), col("risk_table_map")("count")("reward_sell").alias("reward_sell") )
步骤3:动态提取列(索引名称不固定场景)
如果index名称是动态变化的,先收集所有可能的索引值,再自动生成对应列:
import org.apache.spark.sql.functions.{explode, map_keys} // 收集所有唯一的index值 val indexValues = dfWithMap .select(explode(map_keys(col("risk_table_map")))) .distinct() .as[String] .collect() // 动态生成每个index对应的risk_buy列 val indexColumns = indexValues.map(index => col("risk_table_map")(index)("risk_buy").alias(index) ) // 组合所有需要的列 val finalDfDynamic = dfWithMap .select( col("symbol_id") +: indexColumns ++ Seq( col("risk_table_map")("count")("reward_buy").alias("reward_buy"), col("risk_table_map")("count")("reward_sell").alias("reward_sell") ): _* )
注意事项
- 若
risk_table存在重复index,map_from_entries会保留最后一个出现的结构体,需提前去重 - 部分行缺失某个
index时,对应列会显示null,可通过coalesce设置默认值 - 需确保Spark版本≥2.4,支持
map_from_entries和transform函数
内容的提问来源于stack exchange,提问作者Eternal Student
相关产品推荐
相关产品推荐

