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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 21:56:01