Spark DataFrame多扁平列转结构体数组的实现方法咨询
Spark DataFrame 扁平列转结构体数组实现
Python 实现
思路
- 识别成对的目标列(如
Fat Value与Fat Measure为一组,对应type为Fat) - 对每组列,用
struct()构造包含type、amount、measure的结构体 - 用
array()将所有结构体组合成数组类型的新列
代码示例
from pyspark.sql import SparkSession from pyspark.sql.functions import struct, array, lit # 初始化SparkSession spark = SparkSession.builder.appName("FlattenToStructArray").getOrCreate() # 构造测试DataFrame(模拟10组Value/Measure列) data = [ (10, "mg", 5, "g", 3, "μg") + (0, "unit")*14 # 补全10组的测试数据 ] columns = [ "Fat Value", "Fat Measure", "Salt Value", "Salt Measure", "Sugar Value", "Sugar Measure", "Protein Value", "Protein Measure", "Fiber Value", "Fiber Measure", "Calcium Value", "Calcium Measure", "Iron Value", "Iron Measure", "Potassium Value", "Potassium Measure", "VitaminA Value", "VitaminA Measure", "VitaminC Value", "VitaminC Measure" ] df = spark.createDataFrame(data, columns) # 构造结构体数组的表达式 structs = [] # 遍历所有列对,提取type名称并构造struct for i in range(0, len(columns), 2): col_value = columns[i] col_measure = columns[i+1] # 从列名中提取type(比如"Fat Value" → "Fat") nutrient_type = col_value.rsplit(" ", 1)[0] structs.append( struct( lit(nutrient_type).alias("type"), df[col_value].alias("amount"), df[col_measure].alias("measure") ) ) # 添加新的数组列 result_df = df.withColumn("nutrients", array(*structs)) # 查看结果 result_df.select("nutrients").show(truncate=False)
关键说明
struct()用于将多个字段组合成结构体,通过别名匹配需求中的type/amount/measurearray()将所有结构体对象打包成数组列- 通过拆分列名提取
type值,适配命名规则一致的列对
Scala 实现
思路
与Python实现逻辑一致:分组列对→构造结构体→组合成数组列
代码示例
import org.apache.spark.sql.SparkSession import org.apache.spark.sql.functions.{array, lit, struct} object FlattenToStructArray { def main(args: Array[String]): Unit = { val spark = SparkSession.builder.appName("FlattenToStructArray").getOrCreate() import spark.implicits._ // 构造测试DataFrame val data = Seq( (10, "mg", 5, "g", 3, "μg") +: Seq.fill(14)(0, "unit"): _* ) val columns = Seq( "Fat Value", "Fat Measure", "Salt Value", "Salt Measure", "Sugar Value", "Sugar Measure", "Protein Value", "Protein Measure", "Fiber Value", "Fiber Measure", "Calcium Value", "Calcium Measure", "Iron Value", "Iron Measure", "Potassium Value", "Potassium Measure", "VitaminA Value", "VitaminA Measure", "VitaminC Value", "VitaminC Measure" ) val df = data.toDF(columns: _*) // 生成结构体数组表达式 val structs = columns.grouped(2).map { case Seq(colValue, colMeasure) => val nutrientType = colValue.split(" ").init.mkString(" ") struct( lit(nutrientType).alias("type"), col(colValue).alias("amount"), col(colMeasure).alias("measure") ) }.toSeq // 添加数组列并查看结果 val resultDf = df.withColumn("nutrients", array(structs: _*)) resultDf.select("nutrients").show(truncate = false) } }
关键说明
- 使用
grouped(2)将列名列表按每2个一组拆分,简化列对的遍历逻辑 col()函数引用DataFrame列,配合struct()构建目标结构体array()接收多个结构体参数,生成最终的数组类型列
内容的提问来源于stack exchange,提问作者soMuchToLearnAndShare
相关产品推荐
相关产品推荐

