FastAPI与Docker环境中multiprocessing进程残留信号量文件致/dev/shm占满的解决方法
FastAPI与Docker环境中multiprocessing进程残留信号量文件致/dev/shm占满的解决方法
我完全懂你现在的困扰——用multiprocessing在FastAPI里做并行处理,功能跑起来没问题,但容器里/dev/shm的信号量文件越堆越多,最后直接把空间占满导致服务崩溃,而且又不能用join()阻塞接口,那并行的意义就没了,确实头疼。
咱们先搞清楚问题根源:你创建的daemon子进程结束后会变成僵尸进程,主进程如果不主动调用join()回收它们,操作系统就不会清理这些进程残留的信号量文件,时间一长自然就把/dev/shm撑爆了。
下面给你几个可行的解决办法,结合你的代码来调整:
方法一:用后台线程自动清理已结束的子进程
这个方法既不会阻塞接口,又能让主进程主动回收僵尸进程,彻底解决信号量文件堆积的问题。
修改你的example_server.py,添加进程跟踪和后台清理逻辑:
import asyncio import threading import time from fastapi import FastAPI from multiprocessing import Queue, Process from example_script import create_thread app = FastAPI( title="Example title", description="Example description" ) names = Queue() # 存储所有活跃的子进程,用锁避免多线程竞争 active_processes = [] process_lock = threading.Lock() def cleanup_finished_processes(): """后台线程:定期检查并清理已结束的子进程""" while True: with process_lock: # 筛选出已经停止的进程 finished_processes = [p for p in active_processes if not p.is_alive()] # 对每个结束的进程调用join,回收系统资源 for p in finished_processes: p.join() active_processes.remove(p) # 每隔10秒检查一次,可根据你的业务调整间隔 time.sleep(10) @app.on_event("startup") async def startup_event(): for i in range(15): names.put(f"process_{i}") # 启动清理线程,设为daemon线程,FastAPI退出时会自动终止 cleanup_thread = threading.Thread(target=cleanup_finished_processes, daemon=True) cleanup_thread.start() @app.post("/example") async def example_endpoint(var_1: int): loop = asyncio.get_running_loop() await loop.run_in_executor(None, lambda: handle_creation_example(var_1)) return {"status": "process started"} def handle_creation_example(var_1: int): # 这里保留你原来的逻辑,创建进程后加入跟踪列表 process = create_thread(var_1, names) with process_lock: active_processes.append(process)
然后调整example_script.py的create_thread函数,让它返回创建的进程对象:
import asyncio import multiprocessing def create_thread(var_1, names): # 假设var_2是你业务逻辑里生成的变量,这里先模拟一下 var_2 = 10 name = names.get() process = multiprocessing.Process(target=do_job, args=[var_1, var_2, names, name], name=name, daemon=True) process.start() # 返回进程对象给主进程跟踪 return process def do_job(var_1, var_2, names, name): """确保子进程无论正常/异常退出,都能完成收尾""" try: asyncio.run(process_data(var_1, var_2)) finally: # 不管成功失败,都把进程名放回队列 names.put(name) # 模拟你的process_data函数 async def process_data(var1, var2): # 这里是你的业务逻辑 await asyncio.sleep(5)
方法二:改用multiprocessing.Manager管理跨进程队列(可选)
你提到担心Manager只能在if __name__ == "__main__"里用,但其实在FastAPI的启动事件里初始化Manager是完全可行的,它能更稳定地处理跨进程的队列共享,避免一些潜在的资源泄漏问题:
修改example_server.py的startup事件:
from multiprocessing import Manager @app.on_event("startup") async def startup_event(): global names # 用Manager创建队列,跨进程更稳定 manager = Manager() names = manager.Queue() for i in range(15): names.put(f"process_{i}") # 启动清理线程...
额外注意事项
- 不要随便手动删除/dev/shm里的sem文件,因为正在运行的进程还在使用它们,强行删除会导致进程异常。
- 确保你的
process_data函数没有未处理的异常,否则子进程可能会异常退出,增加资源泄漏的概率(上面的try-finally已经帮你处理了这种情况)。
这样调整后,后台线程会自动帮你回收已结束的子进程,信号量文件就会被系统自动清理,不会再堆积占满/dev/shm了。
备注:内容来源于stack exchange,提问作者Gábor
相关产品推荐
相关产品推荐

