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

如何从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
  • 若结果显示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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.18 00:05:26