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

如何快速为大文件添加行号?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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.13 17:23:17