如何快速为大文件添加行号?Python ETL程序性能优化求助
优化ETL程序中文件行号添加的并行处理方案
你的核心问题在于CPU密集型任务下Python线程的GIL限制——线程在执行行号拼接、字符串编码这类CPU密集操作时,GIL会被单个线程持有,导致其他线程无法执行,出现阻塞。同时原代码使用全局计数器会引发多线程/多进程下的行号冲突,且逐行IO的效率较低。以下是针对性优化方案:
一、用multiprocessing替代threading
由于行号添加属于CPU密集型任务,multiprocessing能绕过GIL,让每个进程独立利用CPU核心,真正实现并行处理多个文件。每个进程处理单独的文件,无需共享计数器,避免同步开销。
二、优化单文件处理的代码效率
1. 修正行号计数逻辑
原代码使用全局count(0)会导致跨文件行号连续,不符合常规的单文件独立行号需求,改为每个文件单独维护行号计数器。
2. 批量IO操作
逐行读写会频繁触发系统IO调用,改用批量读取+批量写入能大幅降低IO开销,提升处理速度。
3. 简化表头处理逻辑
用更简洁的方式跳过表头,减少分支判断的开销。
三、优化后的代码示例
import multiprocessing def process_single_file(orig_filename, filename_tmp, encoding, header_included): with open(orig_filename, 'rb') as file_in, open(filename_tmp, 'wb') as file_out: current_line_num = 1 # 处理表头 if header_included: # 跳过第一行表头 file_in.readline() # 批量处理行,可根据内存调整batch_size batch_size = 10000 while True: lines = file_in.readlines(batch_size) if not lines: break processed_batch = [] for line in lines: # 直接编码行号前缀,减少字符串格式化开销 line_prefix = f"{current_line_num}\t".encode(encoding) processed_batch.append(line_prefix + line) current_line_num += 1 # 批量写入,减少IO调用 file_out.write(b''.join(processed_batch)) # 计算最后一行的字段数(如果需要) field_len = None if lines: last_line = lines[-1].decode(encoding) field_len = len(last_line.split('\t')) return field_len def run_parallel_processing(file_list, encoding='utf-8', header_included=True, max_workers=10): # 使用进程池并行处理文件 with multiprocessing.Pool(max_workers=max_workers) as pool: task_args = [ (orig_file, f"{orig_file}.tmp", encoding, header_included) for orig_file in file_list ] # 批量执行任务 results = pool.starmap(process_single_file, task_args) # 这里可以添加文件移动等后续操作 # for orig_file, tmp_file in zip(file_list, [f"{f}.tmp" for f in file_list]): # shutil.move(tmp_file, orig_file) return results
四、额外优化建议
- 调整batch_size:根据系统内存情况增大
batch_size(比如10万行),进一步降低IO次数,但注意不要超出内存限制。 - 避免不必要的编码转换:仅在需要计算
field_len时对最后一行做decode,其他操作全程使用字节流。 - 文件移动操作后置:文件移动属于IO密集且原子性操作,放在所有文件处理完成后用单线程执行即可,无需并行。
- 内存映射文件(可选):对于超大文件(4GB+),可以使用
mmap模块将文件映射到内存,减少数据拷贝开销,示例:import mmap with open(orig_filename, 'rb') as f: with mmap.mmap(f.fileno(), length=0, access=mmap.ACCESS_READ) as mm: # 按换行符分割处理,注意处理方式需调整 pass
内容的提问来源于stack exchange,提问作者brillenheini
相关产品推荐
相关产品推荐

