多进程下类属性未保留问题及FastAPI进程清理方案咨询
问题解答:FastAPI中JobScheduler循环对象为None及进程清理方案
为什么cls.loop调用stop()时为None
你使用multiprocessing.Process启动子进程运行_loop方法,但多进程之间的内存空间是完全隔离的。子进程中对cls.loop的赋值仅存在于子进程自己的内存副本里,FastAPI所在的主进程中,cls.loop仍然是初始的None值,因此调用stop时会触发AttributeError。
FastAPI Shutdown时的后台进程清理方案
1. 优先使用异步任务而非多进程
FastAPI基于asyncio运行,直接在主事件循环中创建异步任务是最简洁的方案。这种方式下任务与主进程共享内存,shutdown时可直接取消任务并等待完成,无需跨进程通信:
import asyncio import logging from typing import List from .production_status import Job as ProductionStatusJob class JobScheduler: _tasks: List[asyncio.Task] = [] logger = logging.getLogger("job_scheduler") @classmethod def start(cls) -> None: cls.logger.info("Starting Up (1/2)") loop = asyncio.get_event_loop() cls._tasks.append(loop.create_task(cls._run())) cls.logger.info("Startup Complete (2/2)") @classmethod async def _run(cls) -> None: while True: # 你的任务逻辑 await asyncio.sleep(5) @classmethod async def stop(cls) -> None: cls.logger.info("Shutting Down (1/2)") for task in cls._tasks: task.cancel() try: await task except asyncio.CancelledError: pass cls._tasks.clear() cls.logger.info("Shutdown Complete (2/2)")
2. 若必须用多进程,通过IPC传递停止信号
如果业务逻辑依赖多进程(如CPU密集型任务),不要直接操作子进程的类变量,而是用multiprocessing.Queue或Pipe给子进程发送停止指令,让子进程自行终止循环:
import asyncio import logging from multiprocessing import Process, Queue from .production_status import Job as ProductionStatusJob class JobScheduler: _stop_queue: Queue = None _process: Process = None logger = logging.getLogger("job_scheduler") @classmethod def start(cls) -> None: cls.logger.info("Starting Up (1/2)") cls._stop_queue = Queue() cls._process = Process(target=cls._loop, args=(cls._stop_queue,)) cls._process.start() cls.logger.info("Startup Complete (2/2)") @classmethod def _loop(cls, stop_queue: Queue) -> None: loop = asyncio.get_event_loop() loop.create_task(cls._run(stop_queue)) loop.run_forever() loop.close() @classmethod async def _run(cls, stop_queue: Queue) -> None: while True: if not stop_queue.empty() and stop_queue.get(): break # 你的任务逻辑 await asyncio.sleep(5) asyncio.get_event_loop().stop() @classmethod def stop(cls) -> None: cls.logger.info("Shutting Down (1/2)") cls._stop_queue.put(True) cls._process.join() cls.logger.info("Shutdown Complete (2/2)")
3. 使用第三方库管理后台任务
对于定时或长期运行的任务,可借助成熟的第三方库简化管理:
fastapi-utils:提供RepeatFunction装饰器,快速创建周期性任务,支持shutdown时自动清理apscheduler:功能全面的任务调度库,支持异步执行,自带任务管理与停止机制
4. 监听系统信号停止子进程
在子进程中监听SIGTERM信号(FastAPI关闭时会发送该信号),收到信号后主动停止事件循环:
import asyncio import logging import signal from multiprocessing import Process from .production_status import Job as ProductionStatusJob class JobScheduler: _process: Process = None logger = logging.getLogger("job_scheduler") @classmethod def start(cls) -> None: cls.logger.info("Starting Up (1/2)") cls._process = Process(target=cls._loop) cls._process.start() cls.logger.info("Startup Complete (2/2)") @classmethod def _loop(cls) -> None: loop = asyncio.get_event_loop() # 注册SIGTERM信号处理 signal.signal(signal.SIGTERM, lambda sig, frame: loop.stop()) loop.create_task(cls._run()) loop.run_forever() loop.close() @classmethod async def _run(cls) -> None: while True: # 你的任务逻辑 await asyncio.sleep(5) @classmethod def stop(cls) -> None: cls.logger.info("Shutting Down (1/2)") cls._process.terminate() cls._process.join() cls.logger.info("Shutdown Complete (2/2)")
内容的提问来源于stack exchange,提问作者Andrew Sharifikia
相关产品推荐
相关产品推荐

