登录
推荐 文章 Go 技术 课程 下载 专题 AI
首页 >  文章 >  python教程

如何在 PyArrow 中高效实现按组累计求和(无需转为 Pandas)

时间:2026-08-20 19:33:32 435浏览 收藏

本文介绍在不转换为 Pandas 的前提下,使用原生 PyArrow API 对分组列执行累计求和(cumulative sum per group),核心思路是结合 pc.cumulative_sum、group_by().aggregate() 和 join 实现偏移量校正,适用于大规模数据且性能显著优于纯 Python 循环。

如何在 PyArrow 中高效实现按组累计求和(无需转为 Pandas)

这段内容要解决的是一个很实际的问题:不把数据转成 Pandas,直接用原生 PyArrow API 来完成分组列上的累计求和(cumulative sum per group)。关键做法并不复杂,核心就是把 `pc.cumulative_sum`、`group_by().aggregate()` 和 `join` 结合起来,通过偏移量校正把每个分组内的累计结果对齐。放在大规模数据场景里,这种方式尤其合适,性能也通常会明显好过纯 Python 循环。

想在 PyArrow 里实现按组累计求和,也就是类似 Pandas 里 df.groupby('b')['a'].cumsum() 这样的操作,不能直接照搬高层接口来做。原因很简单:截至 Apache Arrow 15.x,PyArrow 还没有提供可直接支持 cumsum 的分组计算接口。不过,这并不意味着这件事做不了。借助底层计算函数和表操作的组合,依然可以用向量化的方式把这件事高效完成。这里有一个关键前提:分组键(比如 'b' 列)必须是有序的,也就是说,同一组的行需要连续排在一起,例如 ['x','x','x','y','y','y']。如果原始数据本身是乱序的,那就得先用 table.sort_by('b') 做一遍预处理。

以下是完整实现流程:

✅ 核心步骤解析

  1. 全局累计和:对目标列 'a' 直接调用 pc.cumulative_sum(),得到全量累积数组;
  2. 组内总和聚合:使用 table.group_by('b').aggregate([('a', 'sum')]) 获取每组 'a' 的总和(即各组末尾的全局 cumsum 值);
  3. 构造偏移量数组:将各组总和右移一位(首位置补 0),形成每组起始处应减去的“前缀和”;
  4. 关联校正:通过 table.join() 将偏移量映射回原表,再用 pc.subtract() 从全局 cumsum 中减去对应偏移,即得各组独立 cumsum。

? 完整可运行代码

import pyarrow as pa
import pyarrow.compute as pc

# 构造示例数据(注意:'b' 列已按组有序)
table = pa.table({
'a': [1, 2, 3, 4, 5, 6],
'b': ['x', 'x', 'x', 'y', 'y', 'y']
})

# 步骤 1:全局累计和
cs = pc.cumulative_sum(table['a'])

# 步骤 2:按 'b' 分组求和
gs = table.group_by('b').aggregate([('a', 'sum')])

# 步骤 3:构建 offset_sum —— 每组 cumsum 起点的偏移量([0, sum_group0, sum_group0+sum_group1, ...])
offset_sum = pa.concat_arrays([
pa.array([0]),# 第一组起点偏移为 0
gs['a_sum'].chunks[0][:-1]# 后续组偏移 = 前面所有组的 sum(取除最后一项外的所有项)
])

# 步骤 4:关联并校正
offset_table = pa.table({'b': gs['b'], 'offset_sum': offset_sum})
joined = table.join(offset_table, 'b')
a_cumsum = pc.subtract(cs, joined['offset_sum'])

# 输出结果表
result = pa.table({'a_cumsum': a_cumsum, 'b': table['b']})
print(result.to_pandas())

输出:

 a_cumsumb
0 1x
1 3x
2 6x
3 4y
4 9y
515y

⚠️ 注意事项与最佳实践

  • 分组有序性是前提:本方法依赖组内行物理连续。若 table['b'] 无序(如 ['x','y','x','y']),必须先执行 table = table.sort_by('b'),否则结果错误;
  • 内存友好设计:全程使用 Arrow 原生数组与计算函数,避免中间 Pandas 转换,适合 TB 级数据流处理;
  • 性能优势显著:在 10 万行数据测试中,该向量化方案耗时约 5.9 ms,而等效的 Python 循环实现需 579 ms(相差超 97 倍);
  • 扩展性提示:若需其他累积函数(如 cummax, cummin),可类似构造分组边界索引 + pc.take + pc.combine_chunks 实现,但需额外提取分组起止位置(可通过 gsgroup_indicespc.equal 辅助判断)。

该方案体现了 PyArrow “组合式向量化计算”的设计哲学:不追求单一高阶 API,而是通过灵活拼接基础算子,在保持零拷贝与类型安全的前提下,达成媲美 Pandas 的表达力与远超其的性能表现。

相关阅读
更多>
最新阅读
更多>
课程推荐
更多>