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

Python多进程写入CSV/TXT数据丢失问题求解决方案

多进程下CSV/TXT无丢失写入方案

多进程直接写文件出现数据丢失,核心原因是多个进程同时操作同一文件时的竞争冲突——操作系统的文件缓冲区无法处理并发写入,导致数据被覆盖或截断。以下是不用数据库的可行方案:

1. 队列+单独写进程(推荐)

用multiprocessing.Queue作为中间缓冲,所有工作进程只负责生成数据并放入队列,由一个单独的进程专门负责从队列取数据写入文件。这样彻底避免多进程写文件的竞争,同时内存压力可控。

示例代码:

import multiprocessing
import csv

def worker(queue, task_data):
    # 模拟处理任务生成结果
    for item in task_data:
        result = (item, item * 2)  # 示例结果
        queue.put(result)
    # 每个工作进程结束后放入结束标记
    queue.put(None)

def writer(queue, filename):
    with open(filename, 'w', newline='') as f:
        csv_writer = csv.writer(f)
        csv_writer.writerow(['input', 'output'])  # 写入表头
        done = 0
        total_workers = 2  # 要和实际启动的工作进程数一致
        while done < total_workers:
            data = queue.get()
            if data is None:
                done += 1
                continue
            csv_writer.writerow(data)

if __name__ == '__main__':
    queue = multiprocessing.Queue()
    filename = 'result.csv'
    
    # 启动工作进程
    task1 = [1,2,3,4,5]
    task2 = [6,7,8,9,10]
    p1 = multiprocessing.Process(target=worker, args=(queue, task1))
    p2 = multiprocessing.Process(target=worker, args=(queue, task2))
    p1.start()
    p2.start()
    
    # 启动写进程
    w = multiprocessing.Process(target=writer, args=(queue, filename))
    w.start()
    
    # 等待所有进程结束
    p1.join()
    p2.join()
    w.join()

2. 分进程写临时文件,最后合并

每个工作进程写入自己的临时文件(用进程ID或唯一标识命名),全部进程完成后,由主进程将所有临时文件合并为最终文件。这种方式适合数据量极大的场景,避免队列内存溢出。

示例代码:

import multiprocessing
import csv
import os

def worker(process_id, task_data):
    temp_filename = f'temp_{process_id}.csv'
    with open(temp_filename, 'w', newline='') as f:
        writer = csv.writer(f)
        for item in task_data:
            writer.writerow([item, item*2])

if __name__ == '__main__':
    final_filename = 'result.csv'
    tasks = [[1,2,3], [4,5,6], [7,8,9]]
    processes = []
    
    # 启动工作进程,每个进程写自己的临时文件
    for i, task in enumerate(tasks):
        p = multiprocessing.Process(target=worker, args=(i, task))
        processes.append(p)
        p.start()
    
    # 等待所有工作进程完成
    for p in processes:
        p.join()
    
    # 合并临时文件
    with open(final_filename, 'w', newline='') as final_f:
        writer = csv.writer(final_f)
        writer.writerow(['input', 'output'])  # 写入表头
        for i in range(len(tasks)):
            temp_file = f'temp_{i}.csv'
            with open(temp_file, 'r') as f:
                reader = csv.reader(f)
                for row in reader:
                    writer.writerow(row)
            os.remove(temp_file)  # 删除临时文件

3. 文件锁控制并发写入

如果必须多个进程写同一文件,可以用文件锁保证同一时间只有一个进程写入。注意不同操作系统的锁实现不同:Linux/macOS用fcntl,Windows用msvcrt。

示例代码(Linux/macOS):

import multiprocessing
import csv
import fcntl

def safe_write(filename, row):
    with open(filename, 'a', newline='') as f:
        # 加排他锁
        fcntl.flock(f, fcntl.LOCK_EX)
        writer = csv.writer(f)
        writer.writerow(row)
        # 释放锁(文件关闭时自动释放,这里显式写更清晰)
        fcntl.flock(f, fcntl.LOCK_UN)

def worker(task_data, filename):
    for item in task_data:
        safe_write(filename, [item, item*2])

if __name__ == '__main__':
    filename = 'result.csv'
    # 先写入表头
    with open(filename, 'w', newline='') as f:
        csv.writer(f).writerow(['input', 'output'])
    
    tasks = [[1,2,3], [4,5,6]]
    p1 = multiprocessing.Process(target=worker, args=(tasks[0], filename))
    p2 = multiprocessing.Process(target=worker, args=(tasks[1], filename))
    p1.start()
    p2.start()
    p1.join()
    p2.join()

注意:文件锁方案的性能会比前两种差,因为每次写入都要等待锁,适合数据量不大或写入频率不高的场景。

内容的提问来源于stack exchange,提问作者trung

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.08 22:30:08