GUnicorn重启Worker后multiprocessing.Queue无法工作问题排查
问题背景
启动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,无后续响应。
疑问
- 为何重启Worker后,数据管理进程
s无法收到Worker的消息? - 出现“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 绿色

