如何从Flask服务器向RabbitMQ队列发送消息及排查未入队问题
Flask生产者向RabbitMQ队列发送消息的有效方法及排查方案
核心问题定位
你的代码看似无报错但消息未进入队列,最可能的原因是Flask占用了RabbitMQ的默认端口5672,导致RabbitMQ无法正常监听该端口,或代码实际连接的是Flask服务而非RabbitMQ;此外,交换机与队列的绑定关系缺失/不匹配也会导致消息丢失。
修正后的有效代码
from flask import Flask import pika app = Flask(__name__) @app.route("/create/<message>") def create_bucket(message): credentials = pika.PlainCredentials("guest", "guest") try: # 明确指定RabbitMQ端口(默认5672),避免与Flask端口冲突 connection = pika.BlockingConnection(pika.ConnectionParameters( host="localhost", port=5672, credentials=credentials )) channel = connection.channel() # 若已在管理端创建队列,可注释掉此行,避免属性不匹配问题 # channel.queue_declare(queue="TestQueue", durable=True) # 确保交换机名称、路由键与管理端绑定配置一致 channel.basic_publish( exchange="TestExchange", routing_key="TestQueue", body=message, properties=pika.BasicProperties(delivery_mode=2) # 持久化消息 ) connection.close() return "[x] Message sent %s" % message except pika.exceptions.AMQPConnectionError as e: return f"Failed to connect to RabbitMQ: {str(e)}" if __name__ == "__main__": # 修改Flask端口为非5672的端口,比如5000 app.run(debug=True, port=5000)
关键修正点
- 端口分离:Flask改用5000端口,避免占用RabbitMQ的默认5672端口,确保RabbitMQ能正常监听该端口。
- 连接异常捕获:添加连接错误捕获,直接反馈连接问题,避免无报错但消息未发送的隐性问题。
- 移除冗余队列声明:若已在管理端创建队列,注释掉
queue_declare,避免因属性不匹配导致的异常。
问题排查方案
1. 端口冲突排查
- 执行命令查看5672端口占用情况:
- Linux/macOS:
netstat -an | grep 5672 - Windows:
netstat -ano | findstr :5672
- Linux/macOS:
- 若结果显示Flask进程占用5672,立即修改Flask端口。
2. 连接有效性验证
- 保留代码中的连接异常捕获逻辑,访问接口时若返回连接失败提示,说明RabbitMQ未在指定端口运行,或用户名/密码错误。
3. 交换机与队列绑定检查
- 登录RabbitMQ管理端,进入
Exchanges页面点击TestExchange:- 查看
Bindings列表,确认TestQueue已绑定到该交换机。 - 若为
direct/topic类型交换机,确保绑定的路由键与代码中routing_key完全一致;若为fanout交换机,路由键可忽略。
- 查看
4. 消息流向验证
- 在RabbitMQ管理端的
Exchanges页面,查看TestExchange的Message rates:- 若
publish计数增加但队列无消息,说明绑定关系失效或路由键不匹配。 - 若
publish计数未增加,说明消息未成功发送到交换机,检查代码中交换机名称是否拼写错误,或连接是否正常。
- 若
5. 队列属性一致性检查
- 若保留
queue_declare代码,确保其参数(如durable)与管理端创建的TestQueue属性完全一致,否则会抛出队列属性不匹配的异常。
内容的提问来源于stack exchange,提问作者roqou
相关产品推荐
相关产品推荐

