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

如何用Multiprocessing高效处理超大规模行文件?

低内存处理超大规模文件的多进程优化方案

方案1:用带容量限制的JoinableQueue替代普通Queue

普通Queue没有容量限制,会导致读文件的生产者进程把百万级任务全部塞进队列,既占用大量内存,又会因任务堆积拖慢进程调度。换成JoinableQueue并设置maxsize,当队列满时,put()会自动阻塞,让生产速度和消费速度匹配,完全不用手动批量取行,内存占用极低。

示例代码:

from multiprocessing import Process, JoinableQueue
import requests

def worker(task_queue):
    while True:
        domain = task_queue.get()
        if domain is None:  # 终止信号
            task_queue.task_done()
            break
        # 执行你的requests请求逻辑
        try:
            response = requests.get(f"http://{domain}", timeout=5)
            # 处理响应逻辑...
        except Exception as e:
            # 异常处理...
            pass
        task_queue.task_done()

if __name__ == "__main__":
    PROCESS_NUM = 4  # 进程数根据CPU核数或网络带宽调整
    task_queue = JoinableQueue(maxsize=PROCESS_NUM * 2)  # 队列大小设为进程数的1-2倍即可

    # 启动工作进程
    workers = []
    for _ in range(PROCESS_NUM):
        p = Process(target=worker, args=(task_queue,))
        p.start()
        workers.append(p)

    # 读取文件并分发任务
    with open("file.txt", "r") as f:
        for line in f:
            domain = line.strip()
            if domain:  # 跳过空行
                task_queue.put(domain)

    # 给每个进程发送终止信号
    for _ in range(PROCESS_NUM):
        task_queue.put(None)

    # 等待所有任务完成
    task_queue.join()
    for p in workers:
        p.join()

方案2:用multiprocessing.Pool的迭代方法(更简洁)

Pool的imap()或imap_unordered()是迭代式任务分发,不会一次性加载所有任务到内存,而是按需给空闲进程分配任务,底层自动管理队列,代码更简洁,不用手动维护队列和进程生命周期。

示例代码:

from multiprocessing import Pool
import requests

def process_domain(domain):
    if not domain:
        return
    # 执行你的requests请求逻辑
    try:
        response = requests.get(f"http://{domain}", timeout=5)
        return (domain, response.status_code)
    except Exception as e:
        return (domain, str(e))

if __name__ == "__main__":
    PROCESS_NUM = 4
    with Pool(PROCESS_NUM) as pool:
        with open("file.txt", "r") as f:
            # imap逐行传递任务,内存仅保留当前待处理的少量任务
            results = pool.imap(process_domain, (line.strip() for line in f))
            # 遍历处理结果(按需保留)
            for result in results:
                # 处理每个结果的逻辑...
                pass

原方案变慢的核心原因

普通Queue无容量限制,读文件的速度远快于发请求的消费速度,队列里堆积百万级任务后,不仅占用大量内存,还会导致进程调度开销剧增,最终拖慢整体运行效率。上面两个方案都是通过让生产速度与消费速度同步,从根源避免了任务堆积问题。

内容的提问来源于stack exchange,提问作者MC874

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.05 11:20:56