Spark DataFrame按指定字段分组并转换为嵌套JSON记录
Spark DataFrame分组转换为指定JSON格式解决方案
实现思路
- 先将需要嵌套的字段打包为结构体,并映射为目标JSON的键名,注意类型匹配(如数字转字符串)
- 按指定字段分组,用
collect_list聚合同组的结构体为数组 - 组装完整结构后,用
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
相关产品推荐
相关产品推荐

