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

Spark DataFrame按指定字段分组并转换为嵌套JSON记录

Spark DataFrame分组转换为指定JSON格式解决方案

实现思路

  1. 先将需要嵌套的字段打包为结构体,并映射为目标JSON的键名,注意类型匹配(如数字转字符串)
  2. 按指定字段分组,用collect_list聚合同组的结构体为数组
  3. 组装完整结构后,用to_json转换为目标JSON格式

PySpark 代码实现

from pyspark.sql import SparkSession
from pyspark.sql.functions import struct, collect_list, to_json, col

# 初始化SparkSession
spark = SparkSession.builder.appName("DataFrameToJson").getOrCreate()

# 模拟输入DataFrame(替换为你的实际数据)
data = [("prd_lct", 145, 147, "2024-07-22T05:24:14", 1, 1, 14, 126, "008236686661", "35216")]
columns = ["type", "lctNbr", "itmNbr", "lastUpdatedDate", "lctSeqId", "T7797_PRD_LCT_TYP_CD", "FXT_AIL_ID", "pmyVbuNbr", "upcId", "vndModId"]
df = spark.createDataFrame(data, columns)

# 1. 定义嵌套结构体,重命名字段并转换类型
df_with_structs = df.withColumn(
    "location_struct",
    struct(
        col("lctSeqId"),
        col("T7797_PRD_LCT_TYP_CD").alias("prdLctTypCd"),
        col("FXT_AIL_ID").cast("string").alias("fxtAilId")
    )
).withColumn(
    "item_detail_struct",
    struct(
        col("pmyVbuNbr"),
        col("upcId"),
        col("vndModId")
    )
)

# 2. 按指定字段分组,聚合结构体为数组
grouped_df = df_with_structs.groupBy(
    "type", "lctNbr", "itmNbr", "lastUpdatedDate"
).agg(
    collect_list("location_struct").alias("locations"),
    collect_list("item_detail_struct").alias("itemDetails")
)

# 3. 组装完整结构并转换为JSON
result_df = grouped_df.withColumn("json_output", to_json(struct(
    "type", "lctNbr", "itmNbr", "lastUpdatedDate", "locations", "itemDetails"
)))

# 查看结果
result_df.select("json_output").show(truncate=False)

Scala 代码实现

import org.apache.spark.sql.SparkSession
import org.apache.spark.sql.functions.{struct, collect_list, to_json, col}

object DataFrameToJson {
  def main(args: Array[String]): Unit = {
    val spark = SparkSession.builder.appName("DataFrameToJson").getOrCreate()
    import spark.implicits._

    // 模拟输入DataFrame(替换为你的实际数据)
    val data = Seq(("prd_lct", 145, 147, "2024-07-22T05:24:14", 1, 1, 14, 126, "008236686661", "35216"))
    val df = data.toDF("type", "lctNbr", "itmNbr", "lastUpdatedDate", "lctSeqId", "T7797_PRD_LCT_TYP_CD", "FXT_AIL_ID", "pmyVbuNbr", "upcId", "vndModId")

    // 定义嵌套结构体,重命名字段并转换类型
    val dfWithStructs = df.withColumn(
      "location_struct",
      struct(
        col("lctSeqId"),
        col("T7797_PRD_LCT_TYP_CD").alias("prdLctTypCd"),
        col("FXT_AIL_ID").cast("string").alias("fxtAilId")
      )
    ).withColumn(
      "item_detail_struct",
      struct(
        col("pmyVbuNbr"),
        col("upcId"),
        col("vndModId")
      )
    )

    // 按指定字段分组,聚合结构体为数组
    val groupedDf = dfWithStructs.groupBy("type", "lctNbr", "itmNbr", "lastUpdatedDate")
      .agg(
        collect_list("location_struct").alias("locations"),
        collect_list("item_detail_struct").alias("itemDetails")
      )

    // 组装完整结构并转换为JSON
    val resultDf = groupedDf.withColumn("json_output", to_json(struct(
      col("type"), col("lctNbr"), col("itmNbr"), col("lastUpdatedDate"), col("locations"), col("itemDetails")
    )))

    // 查看结果
    resultDf.select("json_output").show(false)
  }
}

常见报错排查

  • 类型不匹配:目标JSON中fxtAilId为字符串类型,原字段是整数时需用cast("string")转换,否则会导致JSON类型不符
  • 字段别名错误:必须给T7797_PRD_LCT_TYP_CD等字段设置对应别名,否则JSON键名不符合要求
  • 分组键遗漏:严格按type、lctNbr、itmNbr、lastUpdatedDate四个字段分组,否则分组逻辑错误
  • collect_list使用错误:确保传入的是已定义的结构体列,而非单独字段

内容的提问来源于stack exchange,提问作者Teodoro Abarca

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.19 17:05:17