如何在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
相关产品推荐
相关产品推荐

