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

如何在PySpark中创建累积求和列及处理累积聚合?

在PySpark中实现累积求和及背后的处理机制

一、创建累积求和列的正确方式

基础全局累积实现

如果需要全局范围内的累积求和,可以直接使用窗口函数,但要注意这种方式默认会触发全局排序,可能导致所有数据集中到单个分区,仅适合小数据集:

from pyspark.sql import SparkSession
from pyspark.sql.window import Window
from pyspark.sql.functions import col, sum

# 初始化SparkSession
spark = SparkSession.builder.getOrCreate()

# 创建示例数据集
df = spark.createDataFrame([
    (1, "Alice", 100),
    (2, "Bob", 200),
    (3, "Charlie", 150),
    (4, "David", 300),
    (5, "Eve", 250)
], ["ID", "Name", "Value"])

# 定义全局排序的窗口规范,范围从起始行到当前行
window_spec = Window.orderBy("ID").rowsBetween(Window.unboundedPreceding, Window.currentRow)
# 添加累积求和列
df = df.withColumn("CumulativeSum", sum(col("Value")).over(window_spec))

df.show()

超大型数据集优化方案

针对超大型数据集,核心是避免全量数据集中到单个分区,优先通过分组拆分数据并行处理:

from pyspark.sql import SparkSession
from pyspark.sql.window import Window
from pyspark.sql.functions import col, sum

spark = SparkSession.builder.getOrCreate()

# 带分组的示例数据集
df = spark.createDataFrame([
    (1, "A", 100),
    (2, "A", 200),
    (3, "B", 150),
    (4, "B", 300),
    (5, "A", 250)
], ["ID", "Group", "Value"])

# 按Group分区、ID排序的窗口规范
# 每个分组的数据会被分配到不同分区,并行计算累积和
window_spec = Window.partitionBy("Group").orderBy("ID").rowsBetween(Window.unboundedPreceding, Window.currentRow)
df = df.withColumn("GroupCumulativeSum", sum(col("Value")).over(window_spec))

df.show()

如果必须做全局累积,可通过调整spark.sql.shuffle.partitions参数(默认200)增加分区数,缓解单个分区的压力。

二、PySpark处理累积聚合的机制

  1. 数据分区与Shuffle

    • 若窗口指定了partitionBy,Spark会先按分区键对数据进行Shuffle,将同一分组的数据聚集到同一个Executor的分区中;
    • 未指定partitionBy时,会触发全局Shuffle,所有数据被发送到单个分区,这也是小数据集方案不适用于超大型数据的原因。
  2. 分区内排序与计算
    每个分区内的数据会按照orderBy指定的字段排序,之后根据rowsBetween/rangeBetween定义的窗口范围,逐行计算累积聚合值:

    • rowsBetween按物理行数定义窗口范围,适合排序列是唯一标识的场景;
    • rangeBetween按排序列的数值范围定义窗口,适合排序列有连续重复值的场景。
  3. 并行执行逻辑
    带partitionBy的窗口操作会在多个分区上并行执行,每个分区独立完成排序和累积计算,最后合并结果,大幅提升超大型数据集的处理效率。

内容的提问来源于stack exchange,提问作者figs_and_nuts

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.15 21:36:31