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

如何在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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.23 16:17:49