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

GUnicorn重启Worker后multiprocessing.Queue无法工作问题排查

GUnicorn Worker进程通信问题分析与解决方案

问题背景

启动GUnicorn Worker进程后,期望Worker仍能接收来自其他进程的数据。当前尝试使用multiprocessing.Queue实现:在Worker fork前启动数据管理进程,通过两个队列与Worker交互——一个用于Worker请求数据,一个用于接收响应。在post_fork钩子中,Worker发送请求并等待响应,完成后才提供服务。
首次运行正常,但手动终止Worker并由GUnicorn重启后,Worker会卡在post_fork方法,无法收到数据管理进程的响应。

复现示例

配置文件config.py

import logging
import os
import multiprocessing
logging.basicConfig(level=logging.INFO)

bind = "localhost:8080"
workers = 1


def s(req_q: multiprocessing.Queue, resp_q: multiprocessing.Queue):
    while True:
        logging.info("Waiting for messages")
        other_pid = req_q.get()
        logging.info("Got a message from %d", other_pid)
        resp_q.put(os.getpid())


m = multiprocessing.Manager()
q1 = m.Queue()
q2 = m.Queue()


proc = multiprocessing.Process(target=s, args=(q1, q2), daemon=True)
proc.start()


def post_fork(server, worker):
    logging.info("Sending request")
    q1.put(os.getpid())
    logging.info("Request sent")
    other_pid = q2.get()
    logging.info("Got response from %d", other_pid)

应用文件app.py

from flask import Flask
app = Flask(__name__)

启动及重启操作

启动命令:

$ gunicorn -c config.py app:app
INFO:root:Waiting for messages
[2023-01-31 14:20:46 +0800] [24553] [INFO] Starting gunicorn 20.1.0
[2023-01-31 14:20:46 +0800] [24553] [INFO] Listening at: http://127.0.0.1:8080 (24553)
[2023-01-31 14:20:46 +0800] [24553] [INFO] Using worker: sync
[2023-01-31 14:20:46 +0800] [24580] [INFO] Booting worker with pid: 24580
INFO:root:Sending request
INFO:root:Request sent
INFO:root:Got a message from 24580
INFO:root:Waiting for messages
INFO:root:Got response from 24574

首次运行日志显示交互正常。手动终止Worker后重启:

$ kill 24580
[2023-01-31 14:22:40 +0800] [24580] [INFO] Worker exiting (pid: 24580)
Error in atexit._run_exitfuncs:
Traceback (most recent call last):
  File "/usr/lib/python3.6/multiprocessing/util.py", line 319, in _exit_function
    p.join()
  File "/usr/lib/python3.6/multiprocessing/process.py", line 122, in join
    assert self._parent_pid == os.getpid(), 'can only join a child process'
AssertionError: can only join a child process
[2023-01-31 14:22:40 +0800] [24553] [WARNING] Worker with pid 24574 was terminated due to signal 15
[2023-01-31 14:22:40 +0800] [29497] [INFO] Booting worker with pid: 29497
INFO:root:Sending request
INFO:root:Request sent

此时Worker卡在post_fork,无后续响应。

疑问

  1. 为何重启Worker后,数据管理进程s无法收到Worker的消息?
  2. 出现“can only join a child process”错误的原因是什么?是否与队列通信问题相关?

环境

  • Python: 3.8.0
  • GUnicorn: 20.1.0
  • OS: Ubuntu 18.04

补充说明

已参考相关问题尝试multiprocessing.Manager.Queue但未解决;因数据不可序列化无法使用HTTP/gRPC,因对象fork时会报错无法用threading.Thread替代进程。


问题分析与解答

疑问1:重启Worker后数据管理进程收不到消息的原因

GUnicorn的Master进程启动时加载配置文件,此时创建的multiprocessing.Manager()和数据管理进程proc属于Master。当Master fork出Worker时,Worker会继承这些队列和进程的引用,但multiprocessing.Manager的底层依赖跨进程通信连接,Worker继承的队列连接在fork后会失效。

首次启动时,Worker是Master直接fork的,队列连接还未失效,所以通信正常。但重启Worker时,新Worker是Master再次fork的,继承的队列连接已经断开,Worker向q1put的数据无法被Manager传递到数据管理进程s,导致s阻塞在req_q.get(),Worker也卡在q2.get()。

疑问2:“can only join a child process”错误原因

这个错误和队列通信直接相关。Worker继承了Master创建的proc对象(数据管理进程的引用),当Worker退出时,Python的atexit钩子会触发multiprocessing的清理逻辑,尝试对所有子进程调用join()。但proc是Master的子进程,并非当前Worker的子进程,Worker调用proc.join()就会触发断言错误。


修复方案

方案1:Master专属初始化+有效队列连接

修改配置文件,确保只有Master进程启动数据管理进程,Worker使用有效的队列连接:

import logging
import os
import multiprocessing
logging.basicConfig(level=logging.INFO)

bind = "localhost:8080"
workers = 1

def s(req_q: multiprocessing.Queue, resp_q: multiprocessing.Queue):
    while True:
        logging.info("Waiting for messages")
        other_pid = req_q.get()
        logging.info("Got a message from %d", other_pid)
        resp_q.put(os.getpid())

# 仅Master进程执行初始化
def master_init():
    global q1, q2, proc
    m = multiprocessing.Manager()
    q1 = m.Queue()
    q2 = m.Queue()
    proc = multiprocessing.Process(target=s, args=(q1, q2), daemon=True)
    proc.start()

# pre_fork钩子确保Master在fork前完成初始化
def pre_fork(server, worker):
    if not hasattr(server, '_master_initialized'):
        master_init()
        server._master_initialized = True

def post_fork(server, worker):
    logging.info("Sending request")
    q1.put(os.getpid())
    logging.info("Request sent")
    other_pid = q2.get()
    logging.info("Got response from %d", other_pid)

方案2:使用原生multiprocessing.Queue

原生队列无需依赖Manager服务端,性能更好,且避免连接失效问题,需确保在Master fork前创建队列和数据管理进程:

import logging
import os
import multiprocessing
logging.basicConfig(level=logging.INFO)

bind = "localhost:8080"
workers = 1

q1 = multiprocessing.Queue()
q2 = multiprocessing.Queue()

def s(req_q: multiprocessing.Queue, resp_q: multiprocessing.Queue):
    while True:
        logging.info("Waiting for messages")
        other_pid = req_q.get()
        logging.info("Got a message from %d", other_pid)
        resp_q.put(os.getpid())

# 仅Master进程启动数据管理进程
if os.getpid() == os.getppid():
    proc = multiprocessing.Process(target=s, args=(q1, q2), daemon=True)
    proc.start()

def post_fork(server, worker):
    logging.info("Sending request")
    q1.put(os.getpid())
    logging.info("Request sent")
    other_pid = q2.get()
    logging.info("Got response from %d", other_pid)

关键注意事项

  • 禁止Worker继承Master的子进程引用,避免退出时触发无效的join()操作。
  • 使用multiprocessing.Manager时,确保Worker使用的队列连接是Master初始化的有效连接。
  • 数据管理进程必须由Master创建,生命周期与Master保持一致,不会随Worker终止而退出。

内容的提问来源于stack exchange,提问作者Green 绿色

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 02:30:59