在PySpark中处理业务报表时,我们经常需要对同一个分组键计算多个指标的聚合值,比如求和、平均值、最大值等。默认情况下,agg操作会把所有聚合结果放在同一行不同的列上,这在某些需要纵向展示指标名称和对应数值的场景下并不友好。本文将详细介绍如何实现多列聚合,并把聚合后的宽表结果转换为行式结构。

一、PySpark多列聚合的基础写法
PySpark的DataFrame API提供了groupBy与agg组合来完成分组聚合。在agg方法中,我们可以同时传入多个列与对应的聚合函数,从而一次性算出若干指标。底层上,Spark SQL会把每个聚合函数映射为一个聚合表达式,并在物理执行计划中生成对应的Aggregate算子,所有指标在同一分区内按分组键合并,最终输出为一行多列。
下面示例构造一个销售数据DataFrame,按地区分组,对销售额与订单数做多列聚合:
from pyspark.sql import SparkSession
from pyspark.sql import functions as F
spark = SparkSession.builder.appName("demo").getOrCreate()
data = [
("华东", 100, 2),
("华东", 200, 3),
("华北", 150, 1),
("华北", 50, 4)
]
df = spark.createDataFrame(data, ["region", "amount", "orders"])
agg_df = df.groupBy("region").agg(
F.sum("amount").alias("total_amount"),
F.avg("amount").alias("avg_amount"),
F.sum("orders").alias("total_orders")
)
agg_df.show()
上述代码运行后,agg_df的 schema 中包含 region、total_amount、avg_amount、total_orders 四列,每一行代表一个地区的聚合结果。这种结构适合宽表存储,但如果要把指标名称与数值并列成两列(即行式),还需要进一步转换。
多列聚合的优点是只需一次 shuffle 就能算出全部指标,效率较高;缺点是结果形态固定为横向,无法直接满足某些BI工具对长表(long format)的要求。
二、使用stack函数将聚合结果转为行式
Spark SQL提供了stack函数,可以把多个列拆成多行。其语法为 stack(n, col1, val1, col2, val2, ...) ,其中 n 表示拆出的列对数,后面交替写列名与列值表达式。我们可以利用它将聚合后的每个指标列展开为(指标名,指标值)的两行结构。
以下代码在刚才的agg_df基础上,用stack把三个指标列转为行式:
row_df = agg_df.select(
"region",
F.expr("stack(3, 'total_amount', total_amount, 'avg_amount', avg_amount, 'total_orders', total_orders) as (metric, value)")
)
row_df.show()
在上面的表达式中,stack(3, ...)表示拆成3组,每组包含指标名和指标值。最终每行记录为某个地区下的一个指标及其数值。这种方式书写紧凑,且只增加一次投影操作,不会引入额外 shuffle,因此在大多数报表场景中性价比最高。
需要注意的是,stack要求所有 value 列类型兼容,若指标中包含字符串与数值混合,应统一转为字符串或双精度。此外,指标名称写在表达式里属于硬编码,如果聚合列很多,维护成本会上升,此时可考虑动态生成 expr 字符串。
三、基于melt思路手动展开多列
有些PySpark版本或团队规范中不鼓励直接写SQL expr,此时可以用unionByName模拟melt(宽转长)过程:对每一个聚合列,单独选出新行,再用union合并。虽然代码量更大,但逻辑清晰,便于调试和类型控制。
示例代码如下:
from pyspark.sql import DataFrame
def melt_agg(agg_df: DataFrame, key_col: str, metric_cols: list):
parts = []
for m in metric_cols:
part = agg_df.select(
key_col,
F.lit(m).alias("metric"),
F.col(m).cast("double").alias("value")
)
parts.append(part)
result = parts[0]
for p in parts[1:]:
result = result.unionByName(p)
return result
metric_cols = ["total_amount", "avg_amount", "total_orders"]
long_df = melt_agg(agg_df, "region", metric_cols)
long_df.show()
这段逻辑把每个指标列变成一个小DataFrame,只包含分组键、固定的指标名和对应的数值,然后纵向拼接。它的优势是类型安全、可读性强,且容易在循环里加入过滤或格式化;劣势是多次select与union会产生更多任务阶段,数据量极大时性能略逊于stack。
实际项目中,若指标数量固定且追求简洁,优先用stack;若指标动态生成或对代码可读性要求高,用melt式展开更稳妥。
四、两种方案对比与选用建议
我们将两种行式转换方式在可读性、性能、灵活性三个维度做简要对比:
| 方案 | 可读性 | 性能 | 灵活性 |
|---|---|---|---|
| stack函数 | 较紧凑,需懂SQL表达式 | 高,无额外shuffle | 中,指标名硬编码 |
| melt手动展开 | 清晰,纯API调用 | 中,多次union | 高,可动态循环 |
从执行原理看,stack在Catalyst优化器里被翻译成同一个Project节点内的表达式展开,而union方式会生成多个逻辑计划再合并,因此前者在Spark UI上的物理计划更简单。对于定期报表任务,建议封装一个根据聚合字典自动产出stack表达式的工具函数,兼顾效率与可维护性。
总之,PySpark多列聚合后转行式展示并不是难题,核心在于理解agg输出形态与stack或union的改写机制。掌握这两种方法,你就能在指标看板、导出文件等场景中自由切换宽表与长表。
PySparkDataFrame多列聚合行式展示修改时间:2026-08-07 01:24:29