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

如何优化Python代码实现500万量级文件快速读取构建DataFrame

500万量级文件读取构建DataFrame性能优化方案

原代码核心瓶颈

  • 全程单线程串行执行,磁盘IO等待时间占了90%以上的耗时,CPU大部分时间处于空等状态
  • 逐次拼接文件路径、打开关闭文件句柄存在大量重复系统调用开销
  • 列表动态append在百万级元素场景下会触发多次全量内存拷贝,额外消耗性能
  • 没有异常兜底逻辑,遇到坏文件、权限不足的文件会直接中断整个流程
  • 没有内存控制逻辑,全量读入文件内容极容易触发OOM导致进程崩溃

具体优化手段

1. 用线程池做并行IO,压满磁盘吞吐

文件读取属于典型的IO密集型任务,Python的GIL会在IO等待时自动释放,用线程池即可实现并行加速,不需要用开销更大的多进程。线程数根据磁盘类型调整即可:机械盘设8-16,SATA固态设16-32,NVMe固态设32-64,线程数过高会因为文件句柄争抢、磁头频繁寻道反而降低性能。

注意:并行读取前先调大系统单进程文件句柄上限,Linux下执行ulimit -n 65535即可,避免报Too many open files错误。

2. 降低遍历与路径处理开销

Python3版本的os.walk默认基于os.scandir实现,遍历返回的文件对象自带路径、文件名属性,不需要手动调用os.path.join拼接路径,能省掉大量字符串拼接的开销。如果有不需要遍历的隐藏目录、系统目录,直接就地修改dirs列表跳过即可,减少无效遍历。

3. 预分配结果内存,减少动态扩容损耗

提前统计总文件数,初始化固定长度的结果列表,通过索引赋值代替动态append,避免列表扩容时的全量内存拷贝,百万级元素下这部分能省15%左右的耗时。

4. 分块处理避免OOM

单文件平均大小按10KB计算,500万文件总大小约50GB,绝大多数机器的内存无法装下全量数据。处理时可以按每1-10万文件为一个批次,读取完一个批次就转成DataFrame存为parquet等二进制临时文件,最后合并所有临时文件即可,不要把所有内容都存在内存里。

优化后参考代码

import os
from concurrent.futures import ThreadPoolExecutor, as_completed
import pandas as pd
from tqdm import tqdm

# 基础配置
TARGET_DIR = "Folder_5M"
# 根据磁盘类型调整线程数
WORKER_COUNT = 32
# 分块大小,每读BATCH_SIZE个文件就存一次临时文件,避免OOM
BATCH_SIZE = 50000

def read_file(file_entry):
    """单个文件读取逻辑,返回(文件名, 文件内容)"""
    try:
        # 直接用Direntry自带的path属性,省去路径拼接开销
        with open(file_entry.path, "rb") as f:
            # 如果后续需要文本内容,直接在这里指定编码解码,省掉二次处理
            # content = f.read().decode("utf-8", errors="ignore")
            content = f.read()
        return (file_entry.name, content)
    except Exception:
        # 遇到权限不足、损坏的文件直接返回空值,不中断整体流程
        return (file_entry.name, None)

if __name__ == "__main__":
    # 第一步:快速遍历收集所有文件对象
    all_files = []
    for root, dirs, files in os.walk(TARGET_DIR):
        # 不需要遍历的目录直接在这里过滤,比如跳过隐藏目录
        # dirs[:] = [d for d in dirs if not d.startswith(".")]
        all_files.extend(files)
    total = len(all_files)
    print(f"扫描完成,共发现{total}个文件,开始读取")

    # 并行读取+分块存储
    with ThreadPoolExecutor(max_workers=WORKER_COUNT) as executor:
        # 按批次提交任务,避免一次性提交几百万个任务占太多内存
        for batch_idx in range(0, total, BATCH_SIZE):
            batch_files = all_files[batch_idx:batch_idx+BATCH_SIZE]
            names = [None]*len(batch_files)
            contents = [None]*len(batch_files)
            futures = [executor.submit(read_file, fe) for fe in batch_files]
            for i, future in enumerate(tqdm(as_completed(futures), total=len(batch_files), desc=f"批次{batch_idx//BATCH_SIZE+1}")):
                name, content = future.result()
                names[i] = name
                contents[i] = content
            # 批次读完直接存临时文件,释放内存
            batch_df = pd.DataFrame({"file_name": names, "content": contents})
            batch_df.to_parquet(f"tmp_batch_{batch_idx//BATCH_SIZE}.parquet", index=False)
            del names, contents, batch_df
    
    # 所有批次处理完后合并成最终DataFrame
    final_df = pd.concat([pd.read_parquet(f) for f in os.listdir(".") if f.startswith("tmp_batch_")], ignore_index=True)
    # 用完记得删除临时文件
    # for f in os.listdir("."):
    #     if f.startswith("tmp_batch_"):
    #         os.remove(f)

额外优化提示

  • 不要用pandas.concat逐行追加构建DataFrame,这种方式比列表收集完一次性转换慢100倍以上
  • 持久化存储优先选parquet格式,比csv读写速度快10倍以上,占用空间仅为csv的1/5左右
  • 机械盘场景不要盲目调高线程数,随机IO寻道开销会随线程数指数上升,反而比单线程更慢
  • Linux环境下可以给磁盘挂载加noatime参数,读文件时不更新访问时间戳,减少不必要的磁盘写入开销

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.28 06:09:20