将requests的response.content放入队列后multiprocessing.Process无法终止
问题分析与解决方案
你的问题根源在于**multiprocessing.Queue内部的 feeder 线程特性**:
当子进程向qout(普通mp.Queue)中放入数据时,Queue会启动一个后台的feeder线程,负责将数据写入到进程间通信的管道中。这个线程是非守护线程——也就是说,即使你的run()函数执行完毕(已经打印了"Infinite loop terminates"),只要这个feeder线程还在等待主进程读取队列中的数据,子进程就会一直保持存活状态,不会终止。
在你的代码里,主进程只完成了qin.join(),但从未从qout中取出任何数据,导致feeder线程一直处于活跃状态,拖住了子进程的退出。
解决方法
最简单的修复方式就是在主进程中读取qout里的所有数据,让feeder线程完成任务后自动退出:
import multiprocessing as mp import queue import requests import time class ChildProcess(mp.Process): def __init__(self, qin, qout): super().__init__() self.qin = qin self.qout = qout self.daemon = True def run(self): while True: try: url = self.qin.get(block=False) r = requests.get(url, verify=False) self.qout.put(r.content) self.qin.task_done() except queue.Empty: break except requests.exceptions.RequestException as e: print(self.name, e) self.qin.task_done() print("Infinite loop terminates") if __name__ == '__main__': qin = mp.JoinableQueue() qout = mp.Queue() for _ in range(5): qin.put('http://en.wikipedia.org') w = ChildProcess(qin, qout) w.start() qin.join() # 新增:读取qout中的所有数据,让feeder线程完成任务 while not qout.empty(): qout.get() time.sleep(1) print(w.name, w.is_alive()) # 现在会输出 False
其他可选方案
改用
JoinableQueue作为输出队列:
如果需要明确等待输出处理完成,可以把qout也换成JoinableQueue,在主进程中调用qout.join(),同时子进程每次put后调用qout.task_done():# 子进程run函数中put后添加task_done self.qout.put(r.content) self.qout.task_done() # 主进程中 qin.join() qout.join()清理requests连接池(可选):
虽然不是这次问题的直接原因,但requests的连接池可能会残留一些后台线程。如果后续遇到类似问题,可以在子进程run()函数末尾手动关闭会话:def run(self): session = requests.Session() try: # 原有逻辑,把requests.get换成session.get r = session.get(url, verify=False) # ... finally: session.close() # 关闭会话,清理连接池线程
内容的提问来源于stack exchange,提问作者oldPadavan
相关产品推荐
相关产品推荐

