如何在Pandas DataFrame中复现Spark collect_list窗口函数效果
Pandas实现按分区排序累积收集列表的方法
完全可以实现和你描述的PySpark窗口逻辑完全一致的效果,下面是简洁的实现方案:
实现步骤
1. 环境准备与测试数据构造
import pandas as pd # 构造和示例一致的初始DataFrame data = [ [1, 1, 10], [1, 2, 11], [1, 2, 12], [3, 1, 13], [2, 1, 14], [2, 1, 15], [2, 1, 16], [4, 1, 17], [4, 2, 18], [4, 3, 19], [4, 4, 19], [4, 5, 20], [4, 5, 20] ] df = pd.DataFrame(data, columns=['A', 'B', 'C']) # 先按A、B排序,对齐Spark的排序逻辑 df = df.sort_values(['A', 'B'], ignore_index=True)
2. 简洁实现方案(小数据量适用)
直接分组后按B值匹配累积C列表,和Spark的range between unbounded preceding and current row逻辑完全一致,相同B值的行会得到相同的累积列表:
df['D'] = df.groupby('A')['B'].transform( lambda b_ser: b_ser.map( lambda b: df.loc[(df['A'] == b_ser.name) & (df['B'] <= b), 'C'].tolist() ) )
3. 高性能实现方案(大数据量适用)
先预计算每个(A,B)组合对应的累积C列表,再映射回原表,避免逐行重复计算,性能提升明显:
# 预计算每个A分区下每个B对应的累积C列表 accum_dict = ( df.sort_values('B') .groupby(['A', 'B'])['C'].apply(list) .groupby('A').apply(lambda x: x.cumsum()) .to_dict() ) # 映射到原表生成D列 df['D'] = df.apply(lambda row: accum_dict[(row['A'], row['B'])], axis=1)
结果验证
运行后打印df即可得到和你给出的Spark计算完全一致的结果:
A B C D 0 1 1 10 [10] 1 1 2 11 [10, 11, 12] 2 1 2 12 [10, 11, 12] 3 2 1 14 [14, 15, 16] 4 2 1 15 [14, 15, 16] 5 2 1 16 [14, 15, 16] 6 3 1 13 [13] 7 4 1 17 [17] 8 4 2 18 [17, 18] 9 4 3 19 [17, 18, 19] 10 4 4 19 [17, 18, 19, 19] 11 4 5 20 [17, 18, 19, 19, 20, 20] 12 4 5 20 [17, 18, 19, 19, 20, 20]
内容的提问来源于stack exchange,提问作者bigdataadd
相关产品推荐
相关产品推荐

