如何在FastAPI中正确终止multiprocessing进程?
问题描述
我有一个FastAPI应用,使用multiprocessing.Process模块启动了一个进程,该进程监听Redis pubsub频道推送的事件。相关代码如下:
订阅函数:
def subscribe_stream(): """ Listens for event pushed through the Redis pubsub channel """ while True: redis = connect_to_redis() event_stream = redis.pubsub() event_stream.subscribe("event_stream") for result in event_stream.listen(): # DO ACTION
FastAPI启动函数:
event_stream = multiprocessing.Process(target=subscribe_stream) api = FastAPI( on_startup=[event_stream.start], on_shutdown=[event_stream.terminate], )
问题:
当我停止FastAPI项目时,event_stream.terminate会被调用,但由于While True循环,该命令似乎被忽略。请问有没有办法让订阅流保持开启,直到FastAPI触发on_shutdown时再终止?
解决方案
方案1:使用multiprocessing.Event实现优雅退出
用事件标志通知子进程主动停止,替代强制终止,让进程可以自行清理资源后退出。
修改后代码:
import multiprocessing from fastapi import FastAPI import redis def connect_to_redis(): # 补充你的Redis连接逻辑 return redis.Redis(host="localhost", port=6379, db=0) def subscribe_stream(stop_event: multiprocessing.Event): """监听Redis pubsub频道,收到停止信号后退出""" redis_conn = connect_to_redis() event_stream = redis_conn.pubsub() event_stream.subscribe("event_stream") # 用带超时的get_message替代阻塞的listen,方便检查退出标志 while not stop_event.is_set(): message = event_stream.get_message(timeout=1) if message and message["type"] == "message": # 执行你的业务逻辑 print(f"收到消息: {message['data']}") # 清理资源 event_stream.unsubscribe() redis_conn.close() # 创建退出事件 stop_event = multiprocessing.Event() event_stream = multiprocessing.Process(target=subscribe_stream, args=(stop_event,)) api = FastAPI( on_startup=[event_stream.start], on_shutdown=[lambda: stop_event.set()] )
说明:
- 子进程每次循环都会检查
stop_event状态,FastAPI触发shutdown时设置该事件,进程会在下一次循环时退出 get_message(timeout=1)避免无限阻塞,确保退出信号能被及时响应
方案2:捕获终止信号实现退出
如果需要保留terminate()调用,可以在子进程中捕获SIGTERM信号,触发退出逻辑。
修改后代码:
import signal import multiprocessing from fastapi import FastAPI import redis def connect_to_redis(): return redis.Redis(host="localhost", port=6379, db=0) def handle_exit(signum, frame): global should_exit should_exit = True def subscribe_stream(): global should_exit should_exit = False # 注册信号处理函数,捕获terminate()发送的SIGTERM信号 signal.signal(signal.SIGTERM, handle_exit) redis_conn = connect_to_redis() event_stream = redis_conn.pubsub() event_stream.subscribe("event_stream") for message in event_stream.listen(): if should_exit: break if message and message["type"] == "message": # 执行你的业务逻辑 print(f"收到消息: {message['data']}") # 清理资源 event_stream.unsubscribe() redis_conn.close() event_stream = multiprocessing.Process(target=subscribe_stream) api = FastAPI( on_startup=[event_stream.start], on_shutdown=[event_stream.terminate], )
说明:
- 子进程收到
terminate()的信号后,会设置退出标志,在下一次消息处理时跳出循环 - 注意:
listen()是阻塞的,若没有消息推送,退出逻辑要等到下一条消息到来才会触发,不如方案1可靠
方案3:修复原代码的循环冗余问题
原代码每次while True都会重新创建Redis连接和订阅,既浪费资源也影响退出响应。把连接逻辑移到外层,仅在连接断开时重连:
import multiprocessing from fastapi import FastAPI import redis import time def connect_to_redis(): return redis.Redis(host="localhost", port=6379, db=0) def subscribe_stream(stop_event: multiprocessing.Event): while not stop_event.is_set(): try: redis_conn = connect_to_redis() event_stream = redis_conn.pubsub() event_stream.subscribe("event_stream") for message in event_stream.listen(): if stop_event.is_set(): break if message and message["type"] == "message": # 执行你的业务逻辑 print(f"收到消息: {message['data']}") except redis.ConnectionError: # Redis连接断开后,等待5秒再重连 time.sleep(5) finally: # 确保资源被清理 try: event_stream.unsubscribe() redis_conn.close() except: pass stop_event = multiprocessing.Event() event_stream = multiprocessing.Process(target=subscribe_stream, args=(stop_event,)) api = FastAPI( on_startup=[event_stream.start], on_shutdown=[lambda: stop_event.set()] )
说明:
- 外层循环处理重连逻辑,避免重复创建连接
- 内层监听消息时检查退出标志,确保收到信号后及时退出
- 添加异常处理,防止Redis连接异常导致进程崩溃
内容的提问来源于stack exchange,提问作者michael holstein
相关产品推荐
相关产品推荐

