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

使用Python多进程apply_async写入CSV时出现数据丢失问题

修复多进程CSV写入数据丢失的问题

你碰到的这个数据丢失、空白行问题,百分百是多进程同时写入同一个文件的竞争冲突导致的——多个进程抢着打开output.csv写数据,操作系统没法保证这些写入操作的顺序和原子性,结果就是有的写入被覆盖,或者出现截断的空白行。下面给你几个靠谱的修复方案,按推荐程度排序:

方案1:用进程队列统一管理写入(最推荐)

核心思路是:让所有解析进程只负责爬取和解析数据,把结果放到一个进程安全的队列里;单独开一个专门的进程,负责从队列里取数据并写入文件。这样就彻底避免了多进程直接操作文件的冲突。

修改后的代码如下:

import csv
import requests
import re
from multiprocessing import Pool, Queue, Process

def parse_data(url, line_num, queue):
    print(line_num, url)
    try:
        r = requests.get(url)
        htmltext = r.text.encode("utf-8")
        pois = re.findall(re.compile(r'<pois>(.+?)</pois>'), htmltext)
        for poi in pois:
            # 把解析好的数据放到队列里,而不是直接写入
            queue.put(poi)
    except Exception as e:
        print(f"Error processing line {line_num}: {str(e)}")

def write_worker(queue):
    # 专门的写入进程,持续从队列取数据写入
    with open('output.csv', 'ab') as resfile:
        writer = csv.writer(resfile)
        while True:
            poi = queue.get()
            # 收到结束信号就退出
            if poi is None:
                break
            writer.writerow([poi])

def main():
    # 创建进程安全的队列
    queue = Queue()
    # 启动写入进程
    writer_process = Process(target=write_worker, args=(queue,))
    writer_process.start()

    pool = Pool(processes=4)
    with open("input.csv", "rb") as f:
        reader = csv.reader(f)
        for line_num, line in enumerate(reader):
            url = line[0]
            # 把队列传给每个解析进程
            pool.apply_async(parse_data, args=(url, line_num, queue))
    
    pool.close()
    pool.join()

    # 所有解析任务完成后,给队列发结束信号
    queue.put(None)
    writer_process.join()

if __name__ == "__main__":
    main()

关键改动说明:

  • 新增Queue来传递解析后的数据,队列是多进程安全的,不用担心竞争
  • 新增write_worker进程专门负责写入,避免多个进程同时操作文件
  • 解析完成后给队列发送None信号,通知写入进程退出
  • 给parse_data加了异常捕获,避免单个URL出错导致整个进程挂掉

方案2:用文件锁控制写入(适合简单场景)

如果不想改太多代码,可以给写入操作加文件锁,确保同一时间只有一个进程能写入文件。注意:这个方法在Unix/Linux和Windows下的实现略有不同,下面是跨平台的简化版本:

import csv
import requests
import re
from multiprocessing import Pool
import fcntl  # Unix/Linux用这个
# Windows下替换成 import msvcrt

def parse_data(url, line_num):
    print(line_num, url)
    try:
        r = requests.get(url)
        htmltext = r.text.encode("utf-8")
        pois = re.findall(re.compile(r'<pois>(.+?)</pois>'), htmltext)
        for poi in pois:
            write_data(poi)
    except Exception as e:
        print(f"Error processing line {line_num}: {str(e)}")

def write_data(poi):
    with open('output.csv', 'ab') as resfile:
        # 加锁:Unix/Linux用fcntl,Windows用msvcrt.locking
        fcntl.flock(resfile, fcntl.LOCK_EX)
        # Windows替换成:msvcrt.locking(resfile.fileno(), msvcrt.LK_LOCK, 1)
        writer = csv.writer(resfile)
        writer.writerow([poi])
        # 解锁(with块结束会自动关闭文件,也会释放锁,这里可以省略)
        fcntl.flock(resfile, fcntl.LOCK_UN)

def main():
    pool = Pool(processes=4)
    with open("input.csv", "rb") as f:
        reader = csv.reader(f)
        for line_num, line in enumerate(reader):
            url = line[0]
            pool.apply_async(parse_data, args=(url, line_num))
    pool.close()
    pool.join()

if __name__ == "__main__":
    main()

注意事项:

  • 文件锁会让写入操作变成串行,多进程的写入性能会下降
  • Windows和Unix/Linux的锁API不一样,需要根据运行环境调整代码

方案3:每个进程写临时文件,最后合并(适合大数据量)

如果要处理的URL非常多,数据量很大,队列可能占用过多内存,这时候可以让每个进程写自己的临时文件,最后把所有临时文件合并成output.csv:

import csv
import requests
import re
from multiprocessing import Pool
import os
import tempfile

def parse_data(url, line_num, temp_file):
    print(line_num, url)
    try:
        r = requests.get(url)
        htmltext = r.text.encode("utf-8")
        pois = re.findall(re.compile(r'<pois>(.+?)</pois>'), htmltext)
        with open(temp_file, 'ab') as f:
            writer = csv.writer(f)
            for poi in pois:
                writer.writerow([poi])
    except Exception as e:
        print(f"Error processing line {line_num}: {str(e)}")

def main():
    # 创建临时文件列表,每个进程对应一个
    temp_files = [tempfile.mktemp(suffix='.csv') for _ in range(4)]
    pool = Pool(processes=4, initializer=None, initargs=None)
    
    with open("input.csv", "rb") as f:
        reader = csv.reader(f)
        for line_num, line in enumerate(reader):
            url = line[0]
            # 给每个任务分配一个临时文件(按进程索引取模)
            temp_file = temp_files[line_num % 4]
            pool.apply_async(parse_data, args=(url, line_num, temp_file))
    
    pool.close()
    pool.join()

    # 合并所有临时文件到output.csv
    with open('output.csv', 'wb') as output_f:
        writer = csv.writer(output_f)
        for temp_file in temp_files:
            if os.path.exists(temp_file):
                with open(temp_file, 'rb') as temp_f:
                    reader = csv.reader(temp_f)
                    for row in reader:
                        writer.writerow(row)
                # 删除临时文件
                os.remove(temp_file)

if __name__ == "__main__":
    main()

优点:

  • 每个进程独立写自己的文件,完全没有竞争冲突
  • 适合大数据量场景,避免队列内存溢出

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.12 03:49:57