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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 12:57:03