如何在Pandas DataFrame中实现PySpark风格的窗口函数累计求和
错误原因
你当前的代码是按A、B双列分组后在组内做累计求和,仅会对每个A+B分组内部的C值逐行累加,和PySpark的窗口逻辑不匹配。
PySpark中Window.partitionBy('A').orderBy('B')搭配sum聚合的默认窗口规则是:同A分区内,所有B值小于等于当前行B值的行的C值总和,因此同一个A下B值相同的所有行,得到的累计和D是完全一致的。
正确实现
方法1:分组聚合后映射回原表(可读性更高)
import pandas as pd # 构造测试数据 df = pd.DataFrame({ 'A': [1,1,1,3,2,2,2,4,4,4,4,4,4], 'B': [1,2,2,1,1,1,1,1,2,3,4,5,5], 'C': [10,11,12,13,14,15,16,17,18,19,19,20,20] }) # 计算每个(A,B)组的C值总和,再按A分组做累计求和 ab_cum_sum = df.groupby(['A','B'])['C'].sum().groupby('A').cumsum().reset_index(name='D') # 将累计和匹配回原表的对应(A,B)行 df = df.merge(ab_cum_sum, on=['A','B'], how='left').sort_values(['A','B']).reset_index(drop=True)
方法2:分组transform实现(代码更简洁)
# 先按A、B排序保证顺序和Spark逻辑一致 df = df.sort_values(['A','B']) # 按A分组后,组内先按B聚合求和再做累计和,对齐到原表每一行 df['D'] = df.groupby('A', group_keys=False).apply(lambda x: x.groupby('B')['C'].transform('sum').cumsum())
输出结果验证
运行后得到的df和你给出的PySpark结果完全一致:
| A | B | C | D |
|---|---|---|---|
| 1 | 1 | 10 | 10 |
| 1 | 2 | 11 | 33 |
| 1 | 2 | 12 | 33 |
| 2 | 1 | 14 | 45 |
| 2 | 1 | 15 | 45 |
| 2 | 1 | 16 | 45 |
| 3 | 1 | 13 | 13 |
| 4 | 1 | 17 | 17 |
| 4 | 2 | 18 | 35 |
| 4 | 3 | 19 | 54 |
| 4 | 4 | 19 | 73 |
| 4 | 5 | 20 | 113 |
| 4 | 5 | 20 | 113 |
内容的提问来源于stack exchange,提问作者bigdataadd
相关产品推荐
相关产品推荐

