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

Python3.10多进程触发停止信号后偶现挂起问题

多进程队列偶发挂起问题分析

测试脚本

#!/usr/bin/env python3
import multiprocessing as mp


def main():
    queue = mp.Queue()
    stop = mp.Event()
    workers = []
    n = mp.cpu_count()
    print(f"starting {n} processes")
    for i in range(n):
        p = mp.Process(target=work, args=(i, queue, stop))
        workers.append(p)
        p.start()
    print("getting 1000 items from queue")
    for _ in range(1000):
        queue.get()
    print("poisoning processes")
    stop.set()
    print("joining processes")
    for worker in workers:
        # hangs occassionally if terminate not called
        # worker.terminate()
        worker.join()
    print("closing queue")
    queue.close()
    print("returning")


def work(i, queue, stop):
    while not stop.is_set():
        queue.put("something")
    print(f"exiting process {i}")


if __name__ == "__main__":
    main()

挂起后终止的示例输出

λ ./mp_template.py
starting 16 processes
getting 1000 items from queue
poisoning processes
joining processes
exiting process 1
exiting process 2
exiting process 0
exiting process 10
exiting process 9
exiting process 7
exiting process 11
exiting process 5
exiting process 12
exiting process 6
exiting process 4
exiting process 13
exiting process 3
exiting process 8
exiting process 14
exiting process 15
^CTraceback (most recent call last):
  File "/data/repos/mse-408/./mp_template.py", line 37, in <module>
    main()
  File "/data/repos/mse-408/./mp_template.py", line 24, in main
    worker.join()
  File "/usr/lib/python3.10/multiprocessing/process.py", line 149, in join
    res = self._popen.wait(timeout)
  File "/usr/lib/python3.10/multiprocessing/popen_fork.py", line 43, in wait
    return self.poll(os.WNOHANG if timeout == 0.0 else 0)
  File "/usr/lib/python3.10/multiprocessing/popen_fork.py", line 27, in poll
    pid, sts = os.waitpid(self.pid, flag)
KeyboardInterrupt
Process Process-1:
Traceback (most recent call last):
  File "/usr/lib/python3.10/multiprocessing/process.py", line 317, in _bootstrap
    util._exit_function()
  File "/usr/lib/python3.10/multiprocessing/util.py", line 360, in _exit_function
    _run_finalizers()
  File "/usr/lib/python3.10/multiprocessing/util.py", line 300, in _run_finalizers
    finalizer()
  File "/usr/lib/python3.10/multiprocessing/util.py", line 224, in __call__
    res = self._callback(*self._args, **self._kwargs)
  File "/usr/lib/python3.10/multiprocessing/queues.py", line 199, in _finalize_join
    thread.join()
  File "/usr/lib/python3.10/threading.py", line 1096, in join
    self._wait_for_tstate_lock()
  File "/usr/lib/python3.10/threading.py", line 1116, in _wait_for_tstate_lock
    if lock.acquire(block, timeout):
KeyboardInterrupt

问题

为什么程序会偶尔挂起,有时却正常?如果为每个worker调用worker.terminate(),程序总能正常退出,但worker在stop.set()调用后应该会返回——为什么还会出现挂起?


原因分析

核心问题:mp.Queue的后台线程阻塞

每个mp.Queue内部都运行着一个后台线程,负责将进程内的数据复制到共享内存/管道中(跨进程通信的底层存储)。当worker进程退出时,它会执行清理逻辑:等待这个后台线程完成所有未完成的写入操作,确保所有通过queue.put()提交的数据都被写入到队列底层。

你的代码中存在两个关键矛盾:

  1. Worker进程在stop.set()被触发前,会无限循环往队列里塞数据,而主线程只读取了前1000条就停止读取。这意味着队列中会积压大量未被读取的数据。
  2. 当队列的底层缓冲区被填满时,worker的后台线程会阻塞在写入操作上——因为主线程不再取数据,缓冲区没有空间接收新数据。这种阻塞会导致worker进程无法完成退出前的清理,进而让主线程的worker.join()一直等待,最终程序挂起。

为什么是“偶尔”挂起?

挂起与否取决于stop.set()被调用时,队列缓冲区的填充状态:

  • 如果此时缓冲区还没被填满,worker的后台线程能快速完成剩余写入,worker进程顺利退出,join()正常结束。
  • 如果缓冲区刚好被填满(或后续写入导致填满),后台线程阻塞,worker进程无法退出,就会出现挂起。

terminate()为什么能解决问题?

worker.terminate()是强制杀死worker进程,不会等待它执行退出前的清理逻辑(包括等待后台线程完成)。不管队列里有没有积压数据,进程都会立刻终止,主线程的join()就能顺利完成。但这种方式属于暴力终止,可能导致数据丢失、资源泄漏,不建议作为常规解决方案。


正确解决方案

要让worker进程顺利退出,核心是让队列的后台线程能完成所有写入操作,可以通过以下方式实现:

  1. 读完队列中所有剩余数据
    在stop.set()之后,主线程把队列中积压的所有数据都读取完毕,给缓冲区腾出空间,让worker的后台线程顺利完成写入:
print("poisoning processes")
stop.set()
# 读取队列中剩余的所有数据
while True:
    try:
        queue.get_nowait()
    except mp.queues.Empty:
        break
print("joining processes")
  1. 使用mp.JoinableQueue
    JoinableQueue专门为生产者-消费者场景设计,自带task_done()和join()机制,能更优雅地管理队列任务的完成状态,避免数据积压导致的阻塞。

  2. 给queue.put()设置超时
    在worker的put()操作中添加超时时间,避免无限阻塞在写入操作上:

def work(i, queue, stop):
    while not stop.is_set():
        try:
            queue.put("something", timeout=0.1)
        except mp.queues.Full:
            continue
    print(f"exiting process {i}")

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.16 08:20:27