如何高效将多个大CSV对应数据求和并生成单个Parquet文件?
问题场景与现有痛点
我正在开发类似分形火焰的图像生成程序,流程分为两步:
- 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
相关产品推荐
相关产品推荐

