使用PySpark生成含无重复数组的深度嵌套JSON
解决PySpark生成无重复嵌套JSON的问题
核心问题分析
你遇到的两个关键问题:
- 缺失
MANUFACTURE字段:因为之前的代码没把这个字段加入groupBy的分组键中,导致聚合后丢失该字段。 PARAMETER和DIMENSION重复:直接用collect_list会保留所有重复行,没有做去重处理。
解决方案
根据你的需求,我们可以通过先去重再聚合或者使用自动去重的聚合函数两种方式实现,同时确保MANUFACTURE字段被包含在分组逻辑中。
方式1:先去重再聚合(保留顺序)
如果需要保持参数的原始顺序,先对重复的MANUFACTURE+PRODUCT+PARAMETER+DIMENSION组合去重,再用collect_list生成嵌套数组:
from pyspark.sql import functions as F # 假设你的原始DataFrame名为df # 1. 按核心维度去重,避免重复的参数组合 deduped_df = df.dropDuplicates(["MANUFACTURE", "PRODUCT", "PARAMETER", "DIMENSION"]) # 2. 分组聚合,构建嵌套结构 result_df = deduped_df.groupBy("MANUFACTURE", "PRODUCT")\ .agg( F.collect_list( F.struct( F.col("PARAMETER"), F.col("DIMENSION"), F.col("VALUE") # 根据你的实际字段调整,比如替换成你需要的其他字段 ) ).alias("PARAMETERS") ) # 3. 输出为JSON文件 result_df.write.mode("overwrite").json("/path/to/output")
方式2:使用collect_set自动去重(无需提前去重)
如果不需要严格保留顺序,可以直接用collect_set替代collect_list,它会自动基于整个struct的哈希值去重:
from pyspark.sql import functions as F result_df = df.groupBy("MANUFACTURE", "PRODUCT")\ .agg( F.collect_set( F.struct( F.col("PARAMETER"), F.col("DIMENSION"), F.col("VALUE") ) ).alias("PARAMETERS") ) result_df.write.mode("overwrite").json("/path/to/output")
关键说明
- 必须将
MANUFACTURE加入groupBy参数:这样聚合后该字段会被保留在结果中,符合你的嵌套结构要求。 - 去重逻辑的选择:
dropDuplicates适合需要保留顺序的场景,collect_set更简洁但输出数组是无序的,根据你的实际需求选择。
内容的提问来源于stack exchange,提问作者mhanifmm96
相关产品推荐
相关产品推荐

