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

Python多进程超时终止后主进程挂起问题排查求助

问题分析与解决方案

看起来你遇到的主进程挂起问题,核心原因是主进程中对JoinableQueue的join()调用没有得到正确的完成信号,同时超时处理逻辑里的队列操作和主进程的队列操作存在冲突。下面是具体的问题点和修复方案:

核心问题点

  1. 主进程的join()调用阻塞
    主进程中执行q_jobPart_1.join()和q_jobPart_2.join()后,会一直等待队列中所有任务的task_done()被调用。但在超时场景下,Reporter进程清空队列的操作虽然调用了task_done(),但主进程之前已经向q_jobPart_1添加过一次Poison Pill(None),Reporter又重复添加了一次,导致队列的任务计数混乱,join()无法完成。

  2. JoinableQueue.empty()的非线程安全问题
    在Reporter中使用while not self.jobPart1_queue.empty()来清空队列是不可靠的,因为empty()方法不保证在多进程场景下的准确性,可能会有任务被遗漏,导致join()永远等待。

  3. 多余的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")

关键修改说明

  1. 设置守护进程
    给所有Worker和Reporter进程添加self.daemon = True,这样如果主进程意外退出,子进程会自动终止,避免僵尸进程。

  2. 队列操作添加超时
    在get()方法中添加timeout参数,避免进程无限阻塞在空队列上,比如Worker的next_task = self.task_queue.get(timeout=1)。

  3. 优化Poison Pill发送逻辑
    新增_send_poison_pills方法,先检查队列中是否已有Poison Pill,避免重复发送,同时根据活跃进程数量发送对应数量的终止信号,确保所有工作进程都能收到退出指令。

  4. 主进程的join()添加超时
    主进程对队列的join()操作添加timeout参数,避免在超时场景下无限等待,同时在超时后主动终止剩余的子进程。

  5. 移除多余的终止操作
    去掉主进程中不必要的workerJobReporter.terminate()调用,改为在最后统一检查并终止所有活跃子进程,保证资源清理干净。

通过这些修改,超时场景下主进程不会再挂起,能正常退出,同时所有子进程也能正确终止。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:34:20