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
相关产品推荐
相关产品推荐

