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

运行网站时无法获取RabbitMQ全部数据,仅能得到‘hello’的技术问询

问题分析

你的代码仅完成了RabbitMQ连接的建立,完全缺失消息消费的核心逻辑,而且app.run(debug=True)是阻塞调用,后续的connection.close()永远不会执行。你能拿到'hello'大概率是队列中遗留的测试消息,实际没有持续消费队列全部消息的机制。

解决步骤

1. 补充RabbitMQ消费逻辑

需要声明目标队列、定义消息处理回调,同时因为Flask的app.run会阻塞主线程,必须把RabbitMQ消费放到后台线程运行,避免影响Web服务正常启动。

2. 修改后的完整代码示例

import pika
import sys
import os
from threading import Thread
from website import create_app

def consume_rabbitmq_messages():
    # 建立RabbitMQ连接与通道
    credentials = pika.PlainCredentials(username='danny', password='123', erase_on_connect=True)
    parameters = pika.ConnectionParameters(host='10.147.18.152', port=5672, virtual_host='490', credentials=credentials)
    
    try:
        connection = pika.BlockingConnection(parameters)
        channel = connection.channel()
        
        # 声明要消费的队列(需与生产者使用的队列名称完全一致)
        channel.queue_declare(queue='your_target_queue', durable=True)  # 根据实际场景调整参数
        
        # 定义消息处理回调
        def callback(ch, method, properties, body):
            message_content = body.decode()
            print(f"收到消息: {message_content}")
            # 在这里添加你的业务逻辑,比如存储消息到数据库、返回给前端等
            # 处理完成后手动确认消息,避免未处理完成的消息丢失
            ch.basic_ack(delivery_tag=method.delivery_tag)
        
        # 配置消费策略:auto_ack=False启用手动确认
        channel.basic_consume(queue='your_target_queue', on_message_callback=callback, auto_ack=False)
        
        print("开始监听RabbitMQ队列...")
        channel.start_consuming()
    except pika.exceptions.ConnectionClosedByBroker:
        print("RabbitMQ连接被Broker关闭,尝试重连")
        consume_rabbitmq_messages()  # 可选:添加自动重连逻辑
    except pika.exceptions.AMQPChannelError as e:
        print(f"通道错误: {e}")
    except pika.exceptions.AMQPConnectionError as e:
        print(f"连接错误: {e}")

app = create_app()

if __name__ == '__main__':
    # 启动后台消费线程
    rabbitmq_thread = Thread(target=consume_rabbitmq_messages, daemon=True)
    rabbitmq_thread.start()
    
    # 启动Flask服务
    app.run(debug=True)

3. 关键注意事项

  • 替换your_target_queue为实际使用的队列名称,必须和生产者发送消息时的队列名完全匹配。
  • 如果队列是持久化的,queue_declare需添加durable=True,与生产者配置保持一致。
  • auto_ack=False配合手动basic_ack能避免消息丢失,确保消息处理完成后才从队列移除。
  • 后台线程设置daemon=True,保证Flask主进程退出时,消费线程自动终止。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.09 20:25:13