运行网站时无法获取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
相关产品推荐
相关产品推荐

