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

多进程处理大文本文件:如何保持输出与输入记录顺序一致?

大文本文件并行处理后保持输出顺序的解决方案

问题背景

有一个约300GiB的文本文件,包含可变行数的表头和数据记录。采用Python Multiprocessing做并行处理(CPU密集型任务),当前代码可正常运行,但输出记录顺序与输入不一致,需要高效实现顺序一致的方案,同时了解标准高效的实现架构。


示例输入文件 input.txt

# header
# etc... (the number of lines in the header can vary) 
record #1
record #2
record #3
record #4
record #5
record #6
record #7
record #8
record #9
...

示例代码

#!/usr/bin/env python
import sys
import multiprocessing

def worker(queue,lock,process_id):

    while True:
        data = queue.get()
        if data == None:
            break

        # 模拟CPU密集型处理
        for x in range(1000000):
            data.split()

        lock.acquire()
        print( data.rstrip() + " computed by process #" + process_id )
        sys.stdout.flush()
        lock.release()

if __name__ == '__main__':

    queue = multiprocessing.Queue()
    lock = multiprocessing.Lock()
    workers = []
    num_processes = 4

    for process_id in range(num_processes-1):
        p = multiprocessing.Process(target=worker, args=(queue,lock,str(process_id)))
        p.start()
        workers.append(p)

    with open('input.txt') as handler:
        try:
            # 读取表头
            line = next(handler)
            while line.startswith('#'):
                line = next(handler)

            # 向队列发送记录
            while True:
                queue.put(line)
                line = next(handler)

        except StopIteration:
            for p in workers:
                queue.put(None)

        finally:
            for p in workers:
                p.join()

当前输出

record #4 computed by process #2
record #3 computed by process #1
record #1 computed by process #0
record #2 computed by process #3
record #5 computed by process #2
record #6 computed by process #1
record #7 computed by process #0
record #8 computed by process #3
record #9 computed by process #2

一、实现输出顺序与输入一致的高效方案

核心思路是给每条记录绑定序号标记,Worker处理完成后返回(序号、结果),主进程按序号排序输出,或维护有序缓冲区批量输出。

方案1:使用multiprocessing.Pool.imap(轻量场景推荐)

Pool.imap会自动保留输入顺序,底层已封装结果排序逻辑,无需手动管理队列:

#!/usr/bin/env python
import sys
import multiprocessing

def process_line(line):
    # 模拟CPU密集型处理
    for x in range(1000000):
        line.split()
    return line.rstrip() + " computed by process #" + str(multiprocessing.current_process().pid % 4)

if __name__ == '__main__':
    num_processes = 4

    with open('input.txt') as handler:
        # 读取并输出表头
        line = next(handler)
        while line.startswith('#'):
            print(line.rstrip())
            line = next(handler)
        
        # 构造记录迭代器(从第一条非表头记录开始)
        def record_iterator():
            yield line
            for l in handler:
                yield l
        
        # 用Pool.imap保持顺序处理
        with multiprocessing.Pool(num_processes) as pool:
            for result in pool.imap(process_line, record_iterator()):
                print(result)
                sys.stdout.flush()

方案2:带序号的任务队列+结果排序(复杂任务场景)

手动管理进程时,给每条记录加序号,Worker返回(序号、结果),主进程收集后按序号排序输出:

#!/usr/bin/env python
import sys
import multiprocessing
from queue import Empty

def worker(task_queue, result_queue):
    while True:
        try:
            idx, data = task_queue.get(timeout=1)
            if idx is None:
                break
        except Empty:
            continue
        
        # 模拟CPU密集型处理
        for x in range(1000000):
            data.split()
        
        result_queue.put( (idx, data.rstrip() + " computed by process #" + str(multiprocessing.current_process().pid % 4)) )

if __name__ == '__main__':
    num_processes = 4
    task_queue = multiprocessing.Queue(maxsize=num_processes*2)
    result_queue = multiprocessing.Queue()
    workers = []

    # 启动Worker进程
    for _ in range(num_processes):
        p = multiprocessing.Process(target=worker, args=(task_queue, result_queue))
        p.start()
        workers.append(p)

    with open('input.txt') as handler:
        # 读取并输出表头
        line = next(handler)
        while line.startswith('#'):
            print(line.rstrip())
            line = next(handler)
        
        # 发送带序号的任务
        idx = 0
        try:
            task_queue.put( (idx, line) )
            idx +=1
            for line in handler:
                task_queue.put( (idx, line) )
                idx +=1
        except StopIteration:
            pass
        
        # 发送终止信号
        for _ in range(num_processes):
            task_queue.put( (None, None) )
        
        # 收集结果并按序号排序
        results = []
        for _ in range(idx):
            results.append( result_queue.get() )
        
        # 按序号排序后输出
        for res in sorted(results, key=lambda x: x[0]):
            print(res[1])
            sys.stdout.flush()
    
    # 等待Worker进程结束
    for p in workers:
        p.join()

二、更标准高效的大文件并行处理架构

针对300GiB级别的超大型文件,可结合以下优化点提升效率:

  • 分块读取替代逐行读取:按固定大小的块读取文件,Worker内部拆分行(处理跨行边界),减少进程间通信的开销。
  • 避免全局锁与实时输出:让Worker将结果写入各自的临时文件,主进程最后按顺序合并临时文件,完全消除锁竞争。
  • 使用Pool的批量方法:Pool.map/Pool.starmap比手动管理队列更简洁,Python内部已做性能优化;若追求吞吐量,可搭配imap_unordered手动排序。
  • 严格控制内存:保持迭代式处理,避免一次性加载过多数据到内存,确保内存占用稳定。

内容的提问来源于Stack Exchange,提问作者Fravadona

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.26 22:02:37