Python多进程处理时写入文件丢失部分内容,求问题排查
问题:多进程写入文件丢失内容
我尝试读取包含10000行内容的文件,生成归档文件后却只看到9874行,且每次丢失的行数都不固定。脚本执行无报错,但部分迭代内容未写入归档文件,请问我哪里操作有误?
import multiprocessing import hashlib from tqdm import tqdm archive = open('color_archive.txt', 'w') def generate_hash(yellow: str) -> str: b256 = hashlib.sha256(yellow.encode()).hexdigest() x = ' '.join([yellow, b256]) archive.write(f"{x}\n") if __name__ == "__main__": listofcolors = [] with open('x.txt') as f: for yellow in tqdm(f, desc="Generating..."): listofcolors.append(yellow.strip()) cpustotal = cpu_count() - 1 pool = multiprocessing.Pool(cpustotal) results = pool.imap(generate_hash, listofcolors) pool.close() pool.join() print('DONE')
问题原因与解决方法
核心问题:多进程共享文件句柄引发写入冲突
你在主进程中打开archive文件句柄后,直接让多个子进程调用该句柄写文件。多进程环境下,文件指针的状态不会在进程间同步,多个进程同时写入时会出现内容覆盖、指针位置混乱的情况,最终导致部分内容丢失,且丢失行数不固定。解决方法1:子进程返回结果,主进程统一写入(推荐)
让子进程仅负责计算哈希并返回结果,由主进程集中写入文件,从根源避免多进程写冲突:import multiprocessing import hashlib from tqdm import tqdm def generate_hash(yellow: str) -> str: b256 = hashlib.sha256(yellow.encode()).hexdigest() return f"{yellow} {b256}\n" if __name__ == "__main__": listofcolors = [] with open('x.txt') as f: for yellow in tqdm(f, desc="Reading..."): listofcolors.append(yellow.strip()) cpustotal = multiprocessing.cpu_count() - 1 pool = multiprocessing.Pool(cpustotal) # 主进程打开文件,遍历子进程返回的结果写入 with open('color_archive.txt', 'w') as archive: for result in tqdm(pool.imap(generate_hash, listofcolors), total=len(listofcolors), desc="Writing..."): archive.write(result) pool.close() pool.join() print('DONE')解决方法2:使用进程锁控制文件写入
如果必须在子进程中写入文件,需要给写入操作加进程锁,确保同一时间只有一个进程操作文件:import multiprocessing import hashlib from tqdm import tqdm def generate_hash(yellow: str, lock, archive_path): b256 = hashlib.sha256(yellow.encode()).hexdigest() line = f"{yellow} {b256}\n" # 加锁后执行写入操作 with lock: with open(archive_path, 'a') as archive: archive.write(line) if __name__ == "__main__": listofcolors = [] with open('x.txt') as f: for yellow in tqdm(f, desc="Reading..."): listofcolors.append(yellow.strip()) cpustotal = multiprocessing.cpu_count() - 1 # 创建跨进程可用的锁 lock = multiprocessing.Lock() pool = multiprocessing.Pool(cpustotal) # 将锁、文件路径和每行数据一起传给子进程 pool.starmap(generate_hash, [(color, lock, 'color_archive.txt') for color in listofcolors]) pool.close() pool.join() print('DONE')额外修正点
原代码中调用cpu_count()时未加multiprocessing前缀,会引发NameError,上述修复代码已补上完整调用multiprocessing.cpu_count()。
内容的提问来源于stack exchange,提问作者Dasher Love
相关产品推荐
相关产品推荐

