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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.28 05:27:00