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

如何高效将多个大CSV对应数据求和并生成单个Parquet文件?

问题场景与现有痛点

我正在开发类似分形火焰的图像生成程序,流程分为两步:

  1. C++数据生成:生成约7亿个带[R,G,B,A]颜色值的点,归入(9000,9000,4)矩阵的对应bin,将扁平化后的矩阵写入CSV文件,重复生成5-20次。每个CSV含8100万行4列数据,格式如下:
r,g,b,a
3,2,0,4   // 对应bin索引 [0,0]
1,1,2,2   // 对应bin索引 [0,1]
0,0,1,1   // 对应bin索引 [0,2]

生成环境无sudo权限,无法直接写入Parquet格式。
2. Python数据合并与可视化:需要将所有CSV中对应bin的数值求和,转换为(9000,9000,4)结构的numpy数组,最终保存为单个Parquet文件以释放CSV占用的存储空间。

现有代码可实现需求,但存在内存占用过高、运行耗时久的问题——核心原因是一次性加载全量数据到内存,且额外增加了CSV转Parquet的中间冗余步骤。

优化方案

1. 流式读取CSV+实时累加,跳过中间Parquet转存

全程仅保留一个(9000,9000,4)的累加数组,通过分块读取每个CSV文件,实时将对应bin的数值累加到数组中,彻底避免加载全量数据到内存。内存占用仅为累加数组的大小(约1.2GB),远低于原方案的数GB级占用。

代码实现

import numpy as np
import pyarrow as pa
import pyarrow.parquet as pq

# 初始化累加数组:指定uint32 dtype,足够覆盖20次累加的数值范围
accumulator = np.zeros((9000, 9000, 4), dtype=np.uint32)

# 定义所有CSV文件路径(示例为file1.csv到file6.csv)
file_paths = [r"C:\Users\Path\to\file{}.csv".format(i) for i in range(1, 7)]

for path in file_paths:
    print(f"Processing {path}")
    # 分块读取CSV,64MB块大小可根据内存情况调整
    csv_reader = pa.csv.open_csv(
        path,
        read_options=pa.csv.ReadOptions(block_size=64 * 1024 * 1024)
    )
    for batch in csv_reader:
        # 将当前块的颜色列转为numpy数组
        r = batch.column('r').to_numpy()
        g = batch.column('g').to_numpy()
        b = batch.column('b').to_numpy()
        a = batch.column('a').to_numpy()
        
        # 计算当前块每行对应的bin索引:行号 = y*9000 + x
        row_ids = np.arange(len(r))
        x_coords = row_ids % 9000
        y_coords = row_ids // 9000
        
        # 实时累加对应bin的颜色值
        accumulator[y_coords, x_coords, 0] += r
        accumulator[y_coords, x_coords, 1] += g
        accumulator[y_coords, x_coords, 2] += b
        accumulator[y_coords, x_coords, 3] += a

# 将累加数组转换为PyArrow Table,写入Parquet
# 方案1:扁平化后按原CSV列格式存储,兼容后续读取逻辑
flat_accum = accumulator.reshape(-1, 4)
final_table = pa.table({
    'r': flat_accum[:, 0],
    'g': flat_accum[:, 1],
    'b': flat_accum[:, 2],
    'a': flat_accum[:, 3]
})

# 方案2:直接存储三维数组,后续读取可直接转为numpy结构
# final_table = pa.table({'image': pa.array(accumulator)})

# 写入最终Parquet文件
pq.write_table(final_table, r"C:\Users\Path\to\final_combined.parquet")
print("合并完成,已生成最终Parquet文件")

2. 可选:优化C++生成逻辑(若可修改)

如果能调整C代码,确保CSV的行顺序严格对应(0,0)到(8999,8999)的bin索引(即行号与bin索引完全一一对应),可以省略Python中计算x_coords和y_coords的步骤,进一步提升效率。若环境允许,甚至可以让C直接写入二进制数组文件,Python通过mmap直接加载,彻底跳过CSV的IO开销。

3. 细节优化

  • ** dtype选择**:使用uint32而非默认的int64,将累加数组的内存占用减少一半,且完全覆盖20次累加的数值范围。
  • 分块大小调整:根据可用内存调整block_size,比如32MB或128MB,平衡IO效率与内存占用。
  • 减少数组复制:直接使用PyArrow Batch转numpy数组,避免额外的数据复制操作。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.14 01:31:27