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

在PySpark中基于分类列计数值生成分组计数新列

大规模数据集下按Model统计颜色计数并扩展列的解决方案

针对数十亿条级别的数据集,优先采用Spark内置聚合+关联的方案,避免窗口函数带来的额外性能开销,具体实现如下:

方法一:固定颜色列的高效实现

当颜色种类固定(如Red、Blue、Yellow)时,先按Model聚合各颜色的计数,再与原表关联:

from pyspark.sql.functions import count, when, col

# 1. 按Model聚合各颜色的计数
agg_df = df.groupBy("Model").agg(
    count(when(col("Color") == "Red", 1)).alias("Red"),
    count(when(col("Color") == "Blue", 1)).alias("Blue"),
    count(when(col("Color") == "Yellow", 1)).alias("Yellow")
)

# 2. 关联原表,将聚合结果映射到每条记录
result_df = df.join(agg_df, on="Model", how="left")

性能优势

  • 仅触发一次shuffle(聚合阶段),且聚合后的agg_df数据量极小(仅Model维度),可自动触发广播join,避免大规模数据shuffle。
  • 全部使用Spark内置函数,无UDF带来的序列化开销。

方法二:动态颜色列的实现(Pivot)

如果颜色种类不固定,可使用pivot自动生成颜色列,再补全缺失值:

# 1. 按Model分组并Pivot颜色列,统计计数
pivot_df = df.groupBy("Model").pivot("Color").count().fillna(0)

# 2. 关联原表得到结果
result_df = df.join(pivot_df, on="Model", how="left")

适用场景

  • 颜色种类不确定或动态变化时,无需硬编码颜色值,自动适配所有存在的颜色。

关于窗口函数的问题说明

此前使用窗口函数结合size()和collect_set()的方案,会对每个Model分区内的所有记录重复计算聚合结果,对于数十亿条数据,会带来额外的内存占用和shuffle压力。而聚合+关联的方案仅在Model维度计算一次计数,再通过广播关联映射到全量数据,性能差异显著。

内容的提问来源于stack exchange,提问作者jay-elliot

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.19 00:53:32