设置超时后Python multiprocessing仍挂起,请求排查原因
问题排查与修复
核心问题分析
你的代码存在两个致命逻辑错误,导致进程重复返回结果、甚至挂起:
1. 工作进程陷入无限循环重复发送同一结果
_Thread类的_run方法中,拿到输入参数后进入死循环,反复将同一个参数放入结果队列,直到停止事件被触发:
while not self.event_stopped.is_set(): try: self.queue_results.put(arg, timeout=0.05) print(f'{self.name} sent {arg}') except queue.Full: pass
这会导致:
- 进程处理完第一个参数后,无法再读取后续输入任务(比如Process-0拿到
a后,永远卡在这个循环里,不会去读取c) - 结果队列被同一参数填满,主进程会反复收到重复结果
2. 主进程取结果的逻辑与任务分配不匹配
任务分配时,Process-0会收到a和c两个参数,Process-1收到b一个参数,但主进程按i%2轮流从进程结果队列取结果,第三次取结果时会从Process-0的队列读取,但此时Process-0还在反复发送a,根本没处理c,导致主进程无法拿到c的结果,甚至可能挂起。
3. 超时未触发的原因
工作进程一直在重复发送同一结果,结果队列始终有数据,主进程的get操作永远不会触发Empty异常;同时主进程一直在读取结果,结果队列不会满,工作进程的put操作也不会触发Full异常,超时逻辑完全没机会生效。
修复后的代码
import multiprocessing as mp import queue class Multithreading: def __init__(self, n_processes): self._processes = [ _Worker(name=f'Process-{i}') for i in range(n_processes)] def __enter__(self): for process in self._processes: process.start() print(f'Started {process.name}') return self def __exit__(self, exc_type, exc_val, exc_tb): # 给所有进程发送结束信号 for process in self._processes: process.queue_inputs.put(None) # 用None标记任务结束 # 等待所有进程退出 for process in self._processes: process.join() def run(self): args = ['a', 'b', 'c'] n_calls = len(args) # 分发任务 for i, arg in enumerate(args): m = i % len(self._processes) print(f'Sending argument to {self._processes[m].name}') # 确保任务放入队列 while True: try: self._processes[m].queue_inputs.put(arg, timeout=0.05) print(f'Argument {arg} sent to {self._processes[m].name}') break except queue.Full: pass print('All arguments sent') # 收集所有结果(从所有进程的结果队列中取,直到拿到n_calls个结果) results = [] received_count = 0 process_queues = {p.name: p.queue_results for p in self._processes} while received_count < n_calls: for name, q in process_queues.items(): try: res = q.get(timeout=0.05) print(f'Received {res} from {name}') results.append(res) received_count += 1 if received_count >= n_calls: break except queue.Empty: continue print(f'All results collected: {results}') class _Worker(mp.Process): def __init__(self, name): super().__init__(name=name, target=self._run) self.queue_inputs = mp.Queue() self.queue_results = mp.Queue() def _run(self): print(f'Running {self.name}') while True: try: arg = self.queue_inputs.get(timeout=0.05) # 收到None表示任务结束,退出循环 if arg is None: print(f'{self.name} received stop signal, exiting') break print(f'{self.name} received {arg}') # 处理任务(这里原样返回),发送一次结果即可 while True: try: self.queue_results.put(arg, timeout=0.05) print(f'{self.name} sent {arg}') break # 发送成功后退出,处理下一个任务 except queue.Full: pass except queue.Empty: pass if __name__ == '__main__': # 测试一次即可,原大循环会重复创建进程,易导致资源耗尽 with Multithreading(n_processes=2) as m: m.run()
关键修改点
- 重命名
_Thread为_Worker:避免与线程类混淆,更符合多进程语义 - 修复工作进程逻辑:
- 去掉重复发送同一结果的死循环,发送成功后立即退出,回到外层循环读取下一个任务
- 用
None作为任务结束信号,替代原有的Event,逻辑更清晰
- 优化主进程结果收集逻辑:
- 不再固定轮流读取进程队列,而是遍历所有进程的结果队列,直到收集到所有任务结果
- 解决了任务分配与结果读取不匹配的问题
- 取消不必要的
Event:用任务结束信号替代,减少同步复杂度 - 注释原代码的大循环:原代码的
for _ in range(100000000)会反复创建大量进程,易导致系统资源耗尽,测试时只需运行一次即可
内容的提问来源于stack exchange,提问作者Sergeant Salty
相关产品推荐
相关产品推荐

