如何优化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
相关产品推荐
相关产品推荐

