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

如何在FASTAPI端点中调用已初始化的aio_pika连接进行关闭/修改

解决方法

方法1:使用全局变量

这是最直接的实现方式,把RabbitMQ连接对象声明为全局变量,在启动流程中完成赋值,端点里直接调用即可。

修改后的完整代码:

from fastapi import FastAPI
import asyncio
import aio_pika

app = FastAPI()
# 声明全局连接变量
swarm_connection = None

@app.on_event("startup")
async def startup():
    global swarm_connection
    host = "你的RabbitMQ地址"
    login = "用户名"
    passwd = "密码"  # 避免用pass(Python关键字)作为变量名
    swarm_connection = await aio_pika.connect_robust(
        host=host,
        port=5672,
        login=login,
        password=passwd,
        loop=asyncio.get_event_loop()
    )
    # 创建通道并配置消费
    swarm_channel = await swarm_connection.channel()
    await swarm_channel.set_qos(prefetch_count=1)
    org1_queue = await swarm_channel.declare_queue(
        'org1', 
        auto_delete=False, 
        durable=True, 
        arguments={'x-max-priority':1}
    )
    await org1_queue.consume(solve_problem_test)

@app.get("/close")
async def close_pika():
    global swarm_connection
    if swarm_connection and not swarm_connection.is_closed:
        await swarm_connection.close()  # aio_pika的close是异步方法,必须加await
        return {"status": "连接已关闭"}
    return {"status": "连接已经处于关闭状态"}

注意:原启动函数改为异步写法更符合FastAPI规范,无需手动用ensure_future调度任务。

方法2:使用FastAPI的app.state属性

FastAPI应用实例自带state属性,专门用来存储全局状态,比直接用全局变量更规范,可避免命名空间污染。

代码示例:

from fastapi import FastAPI
import asyncio
import aio_pika

app = FastAPI()

@app.on_event("startup")
async def startup():
    host = "你的RabbitMQ地址"
    login = "用户名"
    passwd = "密码"
    # 将连接存入app.state
    app.state.swarm_connection = await aio_pika.connect_robust(
        host=host,
        port=5672,
        login=login,
        password=passwd,
        loop=asyncio.get_event_loop()
    )
    # 创建通道并配置消费
    swarm_channel = await app.state.swarm_connection.channel()
    await swarm_channel.set_qos(prefetch_count=1)
    org1_queue = await swarm_channel.declare_queue(
        'org1', 
        auto_delete=False, 
        durable=True, 
        arguments={'x-max-priority':1}
    )
    await org1_queue.consume(solve_problem_test)

@app.get("/close")
async def close_pika():
    conn = app.state.swarm_connection
    if conn and not conn.is_closed:
        await conn.close()
        return {"status": "连接已关闭"}
    return {"status": "连接已经处于关闭状态"}

方法3:封装成单例类

如果应用逻辑复杂,或需要在多模块中访问RabbitMQ连接,可将连接封装为单例类,确保全局只有一个连接实例,扩展性更强。

代码示例:

from fastapi import FastAPI
import asyncio
import aio_pika

class RabbitMQConnection:
    _instance = None

    def __new__(cls):
        if cls._instance is None:
            cls._instance = super().__new__(cls)
            cls._instance.connection = None
        return cls._instance

    async def connect(self, host, port, login, password):
        if not self.connection or self.connection.is_closed:
            self.connection = await aio_pika.connect_robust(
                host=host,
                port=port,
                login=login,
                password=password,
                loop=asyncio.get_event_loop()
            )
        return self.connection

    async def close(self):
        if self.connection and not self.connection.is_closed:
            await self.connection.close()
            self.connection = None

app = FastAPI()
rabbitmq_conn = RabbitMQConnection()

@app.on_event("startup")
async def startup():
    host = "你的RabbitMQ地址"
    login = "用户名"
    passwd = "密码"
    await rabbitmq_conn.connect(host, 5672, login, passwd)
    # 创建通道并配置消费
    swarm_channel = await rabbitmq_conn.connection.channel()
    await swarm_channel.set_qos(prefetch_count=1)
    org1_queue = await swarm_channel.declare_queue(
        'org1', 
        auto_delete=False, 
        durable=True, 
        arguments={'x-max-priority':1}
    )
    await org1_queue.consume(solve_problem_test)

@app.get("/close")
async def close_pika():
    await rabbitmq_conn.close()
    return {"status": "连接已关闭"}

重要提示:aio_pika的close()是异步方法,必须添加await才能真正执行关闭操作。

内容的提问来源于stack exchange,提问作者user3431083

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.26 01:09:28