PySpark如何实现多列聚合后转换为行格式输出
PySpark 宽表列求和转长表的最简实现
下面是两种比循环union性能更优、代码更简洁的实现方案,都只需要两次操作即可完成:
方案1:通用兼容版(全版本PySpark可用)
基于聚合+stack函数实现,支持动态适配任意数量的数值列:
from pyspark.sql.functions import sum, expr # 取所有需要计算的列名 target_cols = df.columns # 动态生成聚合规则、unpivot规则 agg_rules = [sum(col).alias(col) for col in target_cols] stack_rule = f"stack({len(target_cols)}, {', '.join([f'\'{col}\', {col}' for col in target_cols])}) as (col, sum)" # 执行计算 result_df = df.agg(*agg_rules).select(expr(stack_rule))
方案2:极简版(PySpark 3.0及以上可用)
直接调用内置melt方法完成宽转长,代码可读性更高:
from pyspark.sql.functions import sum target_cols = df.columns result_df = df.agg(*[sum(col).alias(col) for col in target_cols])\ .melt(value_vars=target_cols, variableName="col", valueName="sum")
两种方案都只对原表做一次扫描聚合,避免了循环union带来的多次表扫描开销,列数越多性能优势越明显。
内容的提问来源于stack exchange,提问作者user43107
相关产品推荐
相关产品推荐

