Flask使用pika消费RabbitMQ消息时回调不执行问题排查求助
问题根因如下:
1. Flask调试模式的双进程冲突
你启动Flask时开启了debug=True,该模式下Flask默认会启动两个进程:一个是代码热重载的监控进程,一个是实际处理Web请求的工作进程。两个进程都会初始化RPIclient实例,同时监听answer队列,RabbitMQ会把worker返回的响应消息轮询分发给两个消费者。如果消息被监控进程消费,处理Web请求的工作进程自然收不到回调,也就不会触发on_response的日志和ACK逻辑,RabbitMQ端只会显示消息已被消费。
2. worker代码缺失核心初始化逻辑
你提供的worker.py代码存在多处缺失:
- 没有创建RabbitMQ连接对象
connection,直接调用connection.channel()会直接报错,无法正常启动 - 没有声明
kaldi_expe交换机、request队列,也没有完成队列和交换机的绑定,如果worker先于server启动,会直接因为队列/交换机不存在抛出异常,即使消息能发送也会因为路由匹配失败丢失。
3. 缺失correlation_id匹配逻辑
官方RPC实现中,on_response回调需要先校验收到消息的correlation_id是否和当前请求的correlation_id一致,你当前代码没有做这个校验,即使收到其他请求的响应也会直接赋值给self.response,高并发场景下会出现响应错乱问题。
4. pika连接非线程安全隐患
pika的BlockingConnection不是线程安全的,虽然你当前设置了threaded=False单线程运行Flask,但后续如果开多线程会直接导致连接异常,回调无法触发。
临时修复方案
- 启动Flask时关闭debug模式,或者设置
use_reloader=False避免双进程:app.run(debug=True, use_reloader=False, threaded=False, host='0.0.0.0') - 补全worker的RabbitMQ初始化逻辑,添加交换机、队列声明和绑定代码,和server端保持一致
- 在
on_response回调中添加correlation_id校验逻辑,只有匹配当前请求的id才赋值给self.response
内容的提问来源于stack exchange,提问作者Gautier A.
相关产品推荐
相关产品推荐

