如何用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
相关产品推荐
相关产品推荐

