Flask使用异步RPC客户端时出现每两次请求失败一次问题
问题根因
- pika BlockingConnection 线程不安全
pika的BlockingConnection设计上不支持跨线程操作,你当前在Flask请求线程调用send_request发消息,又在单独的后台线程运行process_data_events处理回调,跨线程操作会导致连接状态错乱,回调收到了响应但请求线程读不到的问题本质是连接层状态冲突,偶尔成功只是时序刚好没触发冲突。 - 共享字典访问无锁
你在_on_response回调里写queue字典、在Flask请求线程里读queue字典,两处都没有加锁,Python字典单个操作虽然是原子的,但多线程场景下可能出现CPU缓存未同步的情况,导致请求线程读到的始终是旧的None值。 - Flask debug模式双进程冲突
你开启了debug=True,Flask默认会启动一个重载进程,相当于会初始化两次RpcClient,两个实例的queue是完全独立的,请求会随机落到两个进程上,刚好对应你说的每两次请求就失败一次的现象。
修复方案
1. 禁用Flask debug重载
修改app.run参数,避免双进程重复初始化:
app.run(debug=True, use_reloader=False, threaded=True, host='0.0.0.0') # 生产环境直接关闭debug:app.run(debug=False, threaded=True, host='0.0.0.0')
2. 给所有queue操作加锁
修改RpcClient相关逻辑,读写共享字典时统一加锁:
def _on_response(self, ch, method, props, body): with self.internal_lock: self.queue[props.correlation_id] = body def send_request(self, payload): corr_id = str(uuid.uuid4()) with self.internal_lock: self.queue[corr_id] = None self.channel.basic_publish(exchange='kaldi_expe', routing_key="kaldi_expe.web.request", properties=pika.BasicProperties( reply_to=self.callback_queue, correlation_id=corr_id, ), body=payload) return corr_id
同时修改Flask端点的读逻辑,加锁读取并清理已处理的key避免内存泄漏:
@app.route('/rpc_call/<payload>') def rpc_call(payload): rpc_client = app.config['RPCclient'] corr_id = rpc_client.send_request(payload) while True: with rpc_client.internal_lock: res = rpc_client.queue.get(corr_id) if res is not None: del rpc_client.queue[corr_id] return res time.sleep(0.1)
3. 可选优化
如果后续需要支撑更高并发,建议把pika的BlockingConnection替换为线程安全的实现,或者改用SelectConnection避免跨线程操作连接的问题。
小提示:你代码中
from flask import Flask, jsonif是笔误,正确导入应该是jsonify。
内容的提问来源于stack exchange,提问作者Gautier A.
相关产品推荐
相关产品推荐

