Python多进程超时终止后主进程挂起问题排查求助
看起来你遇到的主进程挂起问题,核心原因是主进程中对JoinableQueue的join()调用没有得到正确的完成信号,同时超时处理逻辑里的队列操作和主进程的队列操作存在冲突。下面是具体的问题点和修复方案:
核心问题点
主进程的
join()调用阻塞
主进程中执行q_jobPart_1.join()和q_jobPart_2.join()后,会一直等待队列中所有任务的task_done()被调用。但在超时场景下,Reporter进程清空队列的操作虽然调用了task_done(),但主进程之前已经向q_jobPart_1添加过一次Poison Pill(None),Reporter又重复添加了一次,导致队列的任务计数混乱,join()无法完成。JoinableQueue.empty()的非线程安全问题
在Reporter中使用while not self.jobPart1_queue.empty()来清空队列是不可靠的,因为empty()方法不保证在多进程场景下的准确性,可能会有任务被遗漏,导致join()永远等待。多余的
terminate()调用
主进程最后调用workerJobReporter.terminate()完全没必要,因为Reporter进程已经通过处理Poison Pill正常退出了,强制终止可能会导致资源泄漏或队列状态异常。
修复后的代码
下面是调整后的代码,重点修改了超时处理逻辑和主进程的等待逻辑:
#!/usr/bin/python3 import sys, time, datetime, os, math, multiprocessing, random TYPE_JOB_PART_1 = '1' TYPE_JOB_PART_2 = '2' class Worker(multiprocessing.Process): def __init__(self, task_queue, result_queue, task_type): multiprocessing.Process.__init__(self) self.task_queue = task_queue self.result_queue = result_queue self.task_type = task_type self.daemon = True # 设置为守护进程,主进程结束时自动退出 def run(self): while True: try: next_task = self.task_queue.get(timeout=1) # 添加超时,避免无限阻塞 except multiprocessing.queues.Empty: continue if next_task is None: print('%s: is exiting: %s ===============' % (self.name, self.task_type)) self.task_queue.task_done() break try: job_response = next_task() self.task_queue.task_done() if self.task_type == TYPE_JOB_PART_1: self.result_queue.put(do_jobPart_2(job_response)) elif self.task_type == TYPE_JOB_PART_2: self.result_queue.put(do_reporting(job_response)) except Exception as e: print(f"Error in {self.task_type} worker: {e}") self.task_queue.task_done() return class Reporter(multiprocessing.Process): def __init__(self, task_queue, result_queue, num_tasks, jobPart1_queue, jobPart2_queue): multiprocessing.Process.__init__(self) self.task_queue = task_queue self.result_queue = result_queue self.jobPart1_queue = jobPart1_queue self.jobPart2_queue = jobPart2_queue self.num_tasks = num_tasks self.time_start = datetime.datetime.now() self.time_wait_to_terminate = 3 self.daemon = True # 设置为守护进程 def run(self): processed = 0 while processed < self.num_tasks: self.time_elapsed = (datetime.datetime.now() - self.time_start).total_seconds() if self.time_elapsed > self.time_wait_to_terminate: print(f"TIME IS UP. {self.time_wait_to_terminate}s elapsed!") # 向所有工作队列发送Poison Pill,无需清空队列(守护进程会自动退出) self._send_poison_pills(self.jobPart1_queue) self._send_poison_pills(self.jobPart2_queue) self._send_poison_pills(self.task_queue) print("TIME IS UP: workers stopped, Reporter shutting down....") break try: next_task = self.task_queue.get(timeout=0.5) except multiprocessing.queues.Empty: continue if next_task is None: self.task_queue.task_done() break try: job_response = next_task() self.task_queue.task_done() self.result_queue.put(job_response) processed += 1 if processed == self.num_tasks: self._send_poison_pills(self.jobPart2_queue) print("JobPart_2 workers will be poisoned") except Exception as e: print(f"Error in Reporter: {e}") self.task_queue.task_done() print('==================================================') print('============ END OF PROCESSING ===================') print('==================================================') return def _send_poison_pills(self, queue): # 避免重复发送Poison Pill,先检查是否已经有None(简单判断) try: # 尝试查看队列头部,不取出 peeked = queue.get(block=False) if peeked is None: queue.put(None) # 放回去 return else: queue.put(peeked) # 非None则放回 except multiprocessing.queues.Empty: pass # 发送对应数量的Poison Pill if queue == self.jobPart1_queue: count = len([w for w in multiprocessing.active_children() if TYPE_JOB_PART_1 in w.name]) elif queue == self.jobPart2_queue: count = len([w for w in multiprocessing.active_children() if TYPE_JOB_PART_2 in w.name]) else: count = 1 for _ in range(count): queue.put(None) class do_reporting(object): def __init__(self, info): self.info = info def __call__(self): try: print(f"{self.info['jobPart1_results']['i']}:do_reporting - is RUNNING ") randtime = 0.5 * random.random() time.sleep(randtime) print( f'jobPart1_time:{self.info["jobPart1_results"]["jobPart1_time"]}, jobPart2_time:{self.info["jobPart2_time"]}, report_time:{randtime}' ) return {'results':self.info,'report_time':randtime} except Exception as e: print(f"error:do_reporting: {e}") class do_jobPart_1(object): def __init__(self, i, t0): self.t0 = t0 self.i = i def __call__(self): try: print(f"{self.i}:do_jobPart_1 - is RUNNING ") randtime = 0.5 * random.random() time.sleep(randtime) time_elapsed = (datetime.datetime.now() - self.t0).total_seconds() return {'i':self.i, 't0':self.t0, 'time_elapsed_job1':time_elapsed, 'jobPart1_time':randtime} except Exception as e: print(f"error:do_jobPart_1: {e}") class do_jobPart_2(object): def __init__(self, info): self.info = info def __call__(self): try: print(f"{self.info['i']}:do_jobPart_2 - is RUNNING ") randtime = 0.5 * random.random() time.sleep(randtime) return {"jobPart1_results":self.info,'jobPart2_time':randtime} except Exception as e: print(f"error:do_jobPart_2: {e}") if __name__ == '__main__': print('==================================================') print('============ START PROCESSING ====================') print('==================================================') # 建立通信队列 q_jobPart_1 = multiprocessing.JoinableQueue() q_jobPart_2 = multiprocessing.JoinableQueue() q_reportTasks = multiprocessing.JoinableQueue() q_results = multiprocessing.Queue() numJobs = 90 numWorkers_jobPart1 = 2 numWorkers_jobPart2 = 2 # 创建工作进程 workersJobPart_1 = [ Worker(q_jobPart_1, q_jobPart_2, TYPE_JOB_PART_1) for i in range(numWorkers_jobPart1) ] workersJobPart_2 = [ Worker(q_jobPart_2, q_reportTasks, TYPE_JOB_PART_2) for i in range(numWorkers_jobPart2) ] workerJobReporter = Reporter(q_reportTasks, q_results, numJobs, q_jobPart_1, q_jobPart_2) # 启动所有进程 print(f"Main PID:{os.getpid()}") for w in workersJobPart_1: w.start() print(f"JobPart_1 PID={w.pid}") for w in workersJobPart_2: w.start() print(f"JobPart_2 PID={w.pid}") workerJobReporter.start() print(f"JobReporter PID={workerJobReporter.pid}") # 添加任务到JobPart1队列 time_start = datetime.datetime.now() for i in range(numJobs): q_jobPart_1.put(do_jobPart_1(i, time_start)) # 等待所有工作进程结束,或者超时 try: # 给JobPart1发送Poison Pill for _ in range(numWorkers_jobPart1): q_jobPart_1.put(None) q_jobPart_1.join(timeout=5) # 添加超时,避免无限等待 q_jobPart_2.join(timeout=5) q_reportTasks.join(timeout=5) except Exception as e: print(f"Join timed out: {e}") # 等待所有子进程退出 for w in workersJobPart_1 + workersJobPart_2: if w.is_alive(): w.terminate() w.join() if workerJobReporter.is_alive(): workerJobReporter.terminate() workerJobReporter.join() print("FINISHED")
关键修改说明
设置守护进程
给所有Worker和Reporter进程添加self.daemon = True,这样如果主进程意外退出,子进程会自动终止,避免僵尸进程。队列操作添加超时
在get()方法中添加timeout参数,避免进程无限阻塞在空队列上,比如Worker的next_task = self.task_queue.get(timeout=1)。优化Poison Pill发送逻辑
新增_send_poison_pills方法,先检查队列中是否已有Poison Pill,避免重复发送,同时根据活跃进程数量发送对应数量的终止信号,确保所有工作进程都能收到退出指令。主进程的
join()添加超时
主进程对队列的join()操作添加timeout参数,避免在超时场景下无限等待,同时在超时后主动终止剩余的子进程。移除多余的终止操作
去掉主进程中不必要的workerJobReporter.terminate()调用,改为在最后统一检查并终止所有活跃子进程,保证资源清理干净。
通过这些修改,超时场景下主进程不会再挂起,能正常退出,同时所有子进程也能正确终止。
内容的提问来源于stack exchange,提问作者Jason

