You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Spark DataFrame多扁平列转结构体数组的实现方法咨询

Spark DataFrame 扁平列转结构体数组实现

Python 实现

思路

  1. 识别成对的目标列(如Fat Value与Fat Measure为一组,对应type为Fat)
  2. 对每组列,用struct()构造包含type、amount、measure的结构体
  3. 用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/measure
  • array()将所有结构体对象打包成数组列
  • 通过拆分列名提取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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.27 16:07:23