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

多个Python进程能否安全并发读取同一CSV文件?

多进程通过Seek读取大CSV的潜在问题与解决方案

这个思路本身是利用多核提升大CSV读取效率的常见尝试,但确实存在几个容易踩坑的地方,我来给你梳理清楚:

1. 最关键的问题:行边界错位

你的代码里直接用seek定位字节偏移,但CSV是行结构化的文件,字节偏移很可能刚好落在某一行的中间位置。举个例子:假设你计算的偏移量指向了某一行的"abc,def"中间的c后面,那csv.reader读到的第一行就会是残缺的",def...",直接导致解析错误,甚至后续所有行的字段都错位。

解决办法

必须先预扫描文件,找到完整行的起始字节位置,再把这些完整的行块分配给各个进程:

  • 可以先打开文件,按一定大小读取缓冲区,在缓冲区里找换行符\n(注意Windows可能是\r\n)的位置,记录每个完整行的起始偏移。
  • 比如先把文件分成N个大致均等的块,然后从每个块的起始位置往后找,直到找到第一个换行符,这个位置才是该进程真正的起始偏移。

2. 磁盘IO的性能瓶颈

虽然多进程并发读取,但磁盘(尤其是机械硬盘)的IO带宽是有限的:

  • 机械硬盘的随机寻道开销很大,如果多个进程频繁seek到不同位置,反而可能比单进程顺序读更慢。
  • SSD的随机IO性能好很多,但过度并发也可能导致IO饱和,反而无法提升效率。

优化建议

  • 如果是机械硬盘,尽量让进程处理的文件块是连续的,减少随机寻道。
  • 可以先把大CSV分割成多个小CSV文件,再分给不同进程处理,这样每个进程都是顺序读,效率更高。

3. 代码里的小细节

你的代码中enumerate(csv.reader(f))从seek后的位置开始读,但如果起始偏移不是行开头,第一行就会出错。另外要注意:

  • 打开文件时尽量指定编码,避免不同环境下的编码问题。
  • 如果CSV有表头,要确保只有一个进程处理表头,或者每个进程都知道跳过表头(如果表头只在文件开头的话)。

可行的改进代码示例

比如先预处理得到每个进程的起始偏移和要读的完整行数:

import csv
import os
from multiprocessing import Pool

def get_line_offsets(filename, chunk_size=1024*1024):
    offsets = [0]
    with open(filename, 'r', encoding='utf-8') as f:
        while True:
            chunk = f.read(chunk_size)
            if not chunk:
                break
            # 找到当前块内所有换行符的位置
            line_breaks = [i for i, c in enumerate(chunk) if c == '\n']
            # 计算每个换行符对应的全局字节偏移
            for pos in line_breaks:
                offsets.append(f.tell() - len(chunk) + pos + 1)
    return offsets

def process_chunk(args):
    filename, start_offset, end_offset = args
    with open(filename, 'r', encoding='utf-8') as f:
        f.seek(start_offset)
        reader = csv.reader(f)
        while f.tell() < end_offset:
            try:
                row = next(reader)
                do_stuff(row[0], row[1], ...)
            except StopIteration:
                break

if __name__ == "__main__":
    filename = "large_file.csv"
    offsets = get_line_offsets(filename)
    # 按CPU核心数分割偏移列表
    num_processes = os.cpu_count()
    chunk_offsets = []
    step = len(offsets) // num_processes
    for i in range(num_processes):
        start = offsets[i*step]
        end = offsets[(i+1)*step] if (i+1)*step < len(offsets) else os.path.getsize(filename)
        chunk_offsets.append((filename, start, end))
    
    with Pool(num_processes) as pool:
        pool.map(process_chunk, chunk_offsets)

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.27 07:22:32