Python多进程同时写入文件异常问题及数字持久化代码修复
解决多进程写入文件内容错乱的问题
问题场景
编写了计算数字持久化的代码,为避免单核心迭代存列表耗尽内存,改为直接写入文件;又因单核心运行速度过慢,改用多进程执行,但多进程同时写入文件时出现内容错乱(如行首重复括号、记录拼接错误),需保证每条[数字, 持久化次数]记录完整写入文件(顺序可后续处理)。
原代码
import time import math from multiprocessing import Pool def fun(num): mult = 1 lis = [int(x) for x in str(num)] for i in lis: mult *= i return mult def fun1(num): count = 0 while len(str(num)) > 1: num = fun(num) count += 1 return count def fun2(num): result = fun1(num) with open("results.txt", "a+") as f: f.write(f"[{num}, {result}]\n") if __name__ == "__main__": upper_bound = 10**2 numbers = range(upper_bound) start = time.time() p = Pool() p.map(fun2, numbers) p.close() p.join() end = time.time() with open("results.txt", "r") as f: things = f.read() print(things) print("The time of execution of above program is :", (end-start) * 10**3, "ms", math.log10(upper_bound))
错误输出示例
[0, 0] [2, 0] [18, 1] [[6, 0] [11, 1] [30, 1] [54, 2] [34, 2] [40, 1] [51, 1] [35, 2] [74, 3] [85, 2] [96, 3] [98, 3] [95, 3]
解决方案
方法1:使用进程锁实现写入互斥
多进程同时写入文件时,操作系统的追加操作虽有原子性,但Python的write调用若包含多字节内容,可能被其他进程打断。通过multiprocessing.Lock保证每次写入操作完整执行。
修改后的代码:
import time import math from multiprocessing import Pool, Lock def fun(num): mult = 1 lis = [int(x) for x in str(num)] for i in lis: mult *= i return mult def fun1(num): count = 0 while len(str(num)) > 1: num = fun(num) count += 1 return count def fun2(args): num, lock = args result = fun1(num) with lock: with open("results.txt", "a+") as f: f.write(f"[{num}, {result}]\n") if __name__ == "__main__": upper_bound = 10**2 numbers = range(upper_bound) start = time.time() lock = Lock() # 将每个数字和锁打包成参数 tasks = [(num, lock) for num in numbers] with Pool() as p: p.map(fun2, tasks) end = time.time() with open("results.txt", "r") as f: things = f.read() print(things) print("The time of execution of above program is :", (end-start) * 10**3, "ms", math.log10(upper_bound))
方法2:用队列让主进程统一写入
让子进程只负责计算结果,将结果存入队列,由主进程单独处理文件写入,彻底避免多进程写文件冲突,同时减少锁的开销。
修改后的代码:
import time import math from multiprocessing import Pool, Queue import threading def fun(num): mult = 1 lis = [int(x) for x in str(num)] for i in lis: mult *= i return mult def fun1(num): count = 0 while len(str(num)) > 1: num = fun(num) count += 1 return count def worker(num, queue): result = fun1(num) queue.put((num, result)) def writer(queue): with open("results.txt", "w") as f: while True: item = queue.get() if item is None: # 结束信号 break num, result = item f.write(f"[{num}, {result}]\n") if __name__ == "__main__": upper_bound = 10**2 numbers = range(upper_bound) start = time.time() queue = Queue() # 启动写入线程 write_thread = threading.Thread(target=writer, args=(queue,)) write_thread.start() with Pool() as p: # 提交所有计算任务 p.starmap(worker, [(num, queue) for num in numbers]) # 发送结束信号 queue.put(None) write_thread.join() end = time.time() with open("results.txt", "r") as f: things = f.read() print(things) print("The time of execution of above program is :", (end-start) * 10**3, "ms", math.log10(upper_bound))
方法3:每个进程写入独立临时文件,最后合并
适合超大规模数据计算场景,每个进程写入自己的临时文件,计算完成后由主进程合并所有临时文件到结果文件,完全避免锁和队列的开销。
核心逻辑示例:
import tempfile import os def fun2(num): result = fun1(num) # 每个进程创建唯一临时文件 with tempfile.NamedTemporaryFile(mode="a", delete=False) as f: f.write(f"[{num}, {result}]\n") return f.name if __name__ == "__main__": upper_bound = 10**2 numbers = range(upper_bound) start = time.time() with Pool() as p: temp_files = p.map(fun2, numbers) # 合并所有临时文件 with open("results.txt", "w") as out_f: for temp_file in temp_files: with open(temp_file, "r") as in_f: out_f.write(in_f.read()) os.unlink(temp_file) end = time.time() # 后续输出逻辑省略
内容的提问来源于stack exchange,提问作者Mark Gyomory
相关产品推荐
相关产品推荐

