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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.04.14 11:58:11