Pyspark动态分组行转列:基于group by/pivot的无UDF实现方案求助
Pyspark 无UDF动态行转列实现
方案1:原生pivot实现(最简洁)
Spark内置pivot函数生成的列名天然符合{group值}_{聚合别名}的格式要求,全程使用原生算子,无UDF性能损耗,适配大数据量场景:
import pyspark.sql.functions as F # 可选:如果group唯一值超过默认1000的限制,先调整参数 # spark.conf.set("spark.sql.pivotMaxValues", 10000) result_df = df.groupBy("id")\ .pivot("group")\ .agg( F.first("A1", ignorenulls=True).alias("A1"), F.first("A2", ignorenulls=True).alias("A2"), F.first("B1", ignorenulls=True).alias("B1"), F.first("B2", ignorenulls=True).alias("B2") )
如果同一个id+group组合下存在多行数据,可根据业务需求将first替换为sum/avg/max等聚合函数。
方案2:手动Case When展开(灵活度更高)
如果需要对列生成逻辑做自定义调整,可以手动拼接列表达式,性能和pivot方案一致:
import pyspark.sql.functions as F # 获取所有动态group值、数值字段列表 vcols = ["A1", "A2", "B1", "B2"] group_vals = [r[0] for r in df.select("group").distinct().collect()] # 生成所有动态列表达式 pivot_cols = [ F.when(F.col("group") == g, F.col(c)).alias(f"{g}_{c}") for g in group_vals for c in vcols ] # 分组聚合取非空值 result_df = df.select("id", *pivot_cols)\ .groupBy("id")\ .agg(*[ F.first(c, ignorenulls=True).alias(c) for c in [f"{g}_{c}" for g in group_vals for c in vcols] ])
可选:输出扁平JSON格式
如果需要最终输出每个id对应一个扁平json对象,可追加以下处理:
result_json_df = result_df.withColumn( "flat_json", F.to_json(F.struct(*result_df.columns[1:])) )
内容的提问来源于stack exchange,提问作者newpyspark
相关产品推荐
相关产品推荐

