PySpark如何按值合并实现嵌套数组扁平化 适配万级动态列场景
PySpark 动态嵌套数组扁平化实现方案
实现思路
- 自动识别所有动态命名的数组列,无需硬编码10000+列名
- 将每个数组列拆分为
id、date、对应值、列名标记的统一长表结构 - 合并所有长表后按
id和date分组转宽表,缺失值自动填充默认值
完整可运行代码
from pyspark.sql import functions as F # --------------- 输入df复现(实际场景可替换为spark.read.json读取所有JSON文件) --------------- df = spark.read.json(sc.parallelize([ """{"id":1,"name_1_a":[{"date":2001,"val":1},{"date":2002,"val":2},{"date":2003,"val":3}],"name_1_b":[{"date":2001,"val":4},{"date":2002,"val":5},{"date":2003,"val":6}],"name_2_a":[{"date":2001,"val":21},{"date":2002,"val":22},{"date":2003,"val":23}],"name_10000_xvz":[{"date":2000,"val":30},{"date":2001,"val":31},{"date":2002,"val":32},{"date":2003,"val":33}]}""", """{"id":2,"name_1_a":[{"date":2001,"val":1},{"date":2002,"val":2},{"date":2003,"val":3}],"name_1_b":[{"date":2001,"val":4},{"date":2002,"val":5},{"date":2003,"val":6}],"name_2_a":[{"date":2001,"val":21},{"date":2002,"val":22},{"date":2003,"val":23}],"name_10000_xvz":[{"date":1990,"val":39},{"date":2000,"val":30},{"date":2001,"val":31},{"date":2002,"val":32},{"date":2003,"val":33},{"date":2004,"val":34}]}}""" ])) # --------------- 核心处理逻辑 --------------- # 自动提取所有动态数组列名,排除固定id列 name_cols = [c for c in df.columns if c != 'id'] # 逐个处理每个数组列,生成统一结构的长表 long_df_list = [] for col_name in name_cols: tmp_df = df.select( "id", F.explode(col_name).alias("struct_item") ).select( "id", F.col("struct_item.date").alias("date"), F.col("struct_item.val").alias("val"), F.lit(col_name).alias("name_col") ) long_df_list.append(tmp_df) # 合并所有长表 all_long_df = long_df_list[0] for tmp in long_df_list[1:]: all_long_df = all_long_df.unionByName(tmp) # 分组转宽表,缺失值补0(若val为字符串类型,将0改为空字符串""即可) result_df = all_long_df.groupBy("id", "date").pivot("name_col").agg(F.first("val")).na.fill(0) # 按id和日期排序输出,和预期结果完全一致 result_df.orderBy("id", "date").show()
方案特性
- 完全适配10000+动态命名的数组列,无需手动编写列名
- 自动兼容不同id、不同列的日期范围差异,缺失值自动填充默认值
- 同时支持int和str类型的val值,仅需调整
na.fill的参数即可
内容的提问来源于stack exchange,提问作者Dan
相关产品推荐
相关产品推荐

