Python百万行大文件处理脚本性能优化咨询
百万行级大文件处理通用指导与代码优化方案
通用优化原则
- 优先批量操作:IO是大文件处理的核心瓶颈,逐行读写会频繁触发用户态/内核态切换,批量操作能大幅减少系统调用次数,利用操作系统缓存提升效率。
- 减少对象频繁创建:避免反复生成小对象(如小DataFrame、临时字符串),减少内存分配与销毁的开销。
- 选择高效IO模式:使用二进制模式(
rb/wb)可减少编码转换开销;内存映射(mmap)适合超大文件,但需注意内存限制。 - 针对性并行处理:CPU密集型任务(如JSON解析)用多进程规避GIL限制;IO密集型任务用多线程提升吞吐量。
现有代码的核心问题
你的代码存在三个关键效率短板:
- 逐行写入错误记录,触发大量IO系统调用;
- 每次生成极小的DataFrame并调用
to_csv,反复初始化CSV写入逻辑; - 未处理行数据索引越界的潜在异常。
优化后的代码实现
方案一:基于Pandas的批量优化
import csv import json import pandas as pd filename = 'filename.csv' BATCH_SIZE = 10000 # 可根据内存调整,建议5000-20000之间 def get_metrics_lines(): with open(filename, "rt") as csvfile: datareader = csv.reader(csvfile, delimiter="|") for row in datareader: # 避免索引越界异常 if len(row) > 8: line = row[8] if '"metrics"' in line: yield line def process_batches(): bad_records = [] metrics_data = [] for line in get_metrics_lines(): try: data = json.loads(line) metrics_data.append(data["metrics"]) except json.JSONDecodeError: bad_records.append(line) # 批量处理有效数据 if len(metrics_data) >= BATCH_SIZE: yield pd.json_normalize(metrics_data) metrics_data = [] # 批量写入错误记录 if len(bad_records) >= BATCH_SIZE: with open('bad_records.txt', 'a') as f: f.write('\n'.join(bad_records) + '\n') bad_records = [] # 处理剩余的最后一批数据 if metrics_data: yield pd.json_normalize(metrics_data) if bad_records: with open('bad_records.txt', 'a') as f: f.write('\n'.join(bad_records) + '\n') # 一次性打开输出文件,批量写入DataFrame with open('name.txt', 'w') as f: first_batch = True for df in process_batches(): # 仅第一批次写入表头 df.to_csv(f, index=False, header=first_batch, sep='\t', escapechar='\\') first_batch = False
方案二:绕过Pandas的极致优化(推荐纯数据转换场景)
如果不需要Pandas的数据分析能力,直接用标准库csv写入可进一步提升速度:
import csv import json filename = 'filename.csv' BATCH_SIZE = 10000 def get_metrics_lines(): with open(filename, "rt") as csvfile: datareader = csv.reader(csvfile, delimiter="|") for row in datareader: if len(row) > 8: line = row[8] if '"metrics"' in line: yield line def process_batches_to_csv(): bad_records = [] metrics_rows = [] field_names = None with open('name.txt', 'w', newline='') as outfile: writer = None for line in get_metrics_lines(): try: data = json.loads(line) metrics = data["metrics"] # 初始化表头(仅第一次执行) if not field_names: field_names = list(metrics.keys()) writer = csv.DictWriter( outfile, fieldnames=field_names, delimiter='\t', escapechar='\\', quoting=csv.QUOTE_MINIMAL ) writer.writeheader() metrics_rows.append(metrics) except json.JSONDecodeError: bad_records.append(line) # 批量写入有效数据 if len(metrics_rows) >= BATCH_SIZE: writer.writerows(metrics_rows) metrics_rows = [] # 批量写入错误记录 if len(bad_records) >= BATCH_SIZE: with open('bad_records.txt', 'a') as f: f.write('\n'.join(bad_records) + '\n') bad_records = [] # 处理剩余数据 if metrics_rows: writer.writerows(metrics_rows) if bad_records: with open('bad_records.txt', 'a') as f: f.write('\n'.join(bad_records) + '\n') process_batches_to_csv()
优化背后的机制
- IO系统调用优化:批量写入将多次小IO合并为一次大IO,减少用户态与内核态的切换次数——每次切换耗时约数百纳秒,百万级数据下累计开销巨大。
- 内存缓存利用:操作系统会对文件IO做页缓存,批量写入能填满缓存页后再一次性刷盘,减少磁盘物理写入次数。
- 对象复用:批量生成DataFrame或直接复用
csv.DictWriter,避免反复初始化对象的元数据与逻辑,降低内存碎片化与GC开销。
内容的提问来源于stack exchange,提问作者Yami Mahō
相关产品推荐
相关产品推荐

