Python多进程写入超大体积CSV/TXT文件出现数据计数不匹配求助
问题根因
- 核心问题是多进程无锁并发追加写入同一个文件,多个进程的写操作会互相覆盖/穿插,导致最终文件行数不稳定
- 代码本身存在语法错误,会引发额外运行异常:
- 变量名不统一:开头定义的块大小是大写
Chunk,后续读取文件时用的是小写chunk,会触发变量未定义报错 - 进程入参格式错误:
args=(each_df)没有加逗号,不会被识别为单元素元组,会报参数数量不匹配错误 - 冗余代码:创建了
mp.Pool对象但全程未使用,属于无效代码
- 变量名不统一:开头定义的块大小是大写
修复方案
推荐优先使用「子进程只负责计算,主进程统一按序写文件」的方案,既避免写冲突,还能保证输出文件的块顺序和原文件一致,不会出现数据乱序问题,实现代码如下:
import pandas as pd import multiprocessing as mp # 路径加r前缀避免转义字符问题 source = r"\\share\usr\data.txt" target = r"\\share\usr\data_masked.txt" chunk = 10000 def process_calc(df): ''' get source df do calc and return newdf ... ''' return newdf if __name__ == '__main__': reader = pd.read_table(source, sep='|', chunksize=chunk, encoding='ANSI') # 用进程池管控进程数量,避免分块过多时创建大量进程导致系统资源耗尽 with mp.Pool(mp.cpu_count()) as pool: # 所有子进程并行计算,返回结果顺序和输入分块顺序完全一致 calc_result = pool.map(process_calc, reader) # 主进程统一按顺序写入文件,完全规避写冲突问题 first_write = True for res_df in calc_result: res_df.to_csv(target, index=None, sep='|', mode='a', header=first_write) first_write = False
如果确实需要子进程自己写文件,可以新增进程锁保证同一时间只有一个进程执行写操作,修改方案如下:
import pandas as pd import multiprocessing as mp from multiprocessing import Lock source = r"\\share\usr\data.txt" target = r"\\share\usr\data_masked.txt" chunk = 10000 # 初始化全局进程锁 lock = Lock() def process_calc(df): ''' get source df do calc and return newdf ... ''' return newdf def calc_frame(df): output_df = process_calc(df) # 写文件前先拿锁,写完自动释放 with lock: output_df.to_csv(target, index=None, sep='|', mode='a', header=False) if __name__ == '__main__': reader = pd.read_table(source, sep='|', chunksize=chunk, encoding='ANSI') with mp.Pool(mp.cpu_count()) as pool: pool.map(calc_frame, reader)
内容的提问来源于stack exchange,提问作者mju
相关产品推荐
相关产品推荐

