如何管理Tornado缓冲区?解决WebSocket内存泄漏问题
Tornado WebSocket慢连接内存泄漏与缓冲区管理方案
核心问题根源
慢连接场景下,后端消息发送速度远高于客户端接收/处理速度,导致Tornado底层IOStream的发送缓冲区持续堆积未送达消息,最终引发内存暴涨。以太网环境因带宽充足,缓冲区能快速排空,因此无此问题。
具体解决方法
1. 基于回调的消息队列控制
不要批量发送所有消息,而是维护待发送队列,仅当前一条消息发送完成后,再发送下一条,从根源避免缓冲区堆积。
示例代码:
class MyWebSocketHandler(tornado.websocket.WebSocketHandler): def open(self): self.pending_messages = [] self.is_sending = False def send_message(self, message): self.pending_messages.append(message) self._process_queue() def _process_queue(self): if not self.is_sending and self.pending_messages: self.is_sending = True msg = self.pending_messages.pop(0) # 利用回调标记消息发送完成 self.write_message(msg, callback=self._on_message_sent) def _on_message_sent(self): self.is_sending = False self._process_queue()
2. 限制缓冲区大小并处理溢出
通过IOStream的set_write_buffer_limits方法设置缓冲区高低水位线,当缓冲区超过上限时触发回调,可选择暂停发送、清理旧消息等操作。
示例代码:
class MyWebSocketHandler(tornado.websocket.WebSocketHandler): def open(self): # 设置缓冲区上限10MB,下限5MB self.ws_connection.stream.set_write_buffer_limits(high=10*1024*1024, low=5*1024*1024) self.ws_connection.stream.on_write_buffer_full = self._on_buffer_full self.send_paused = False def _on_buffer_full(self): self.send_paused = True # 可选:清理旧消息,仅保留最新50条 if len(self.pending_messages) > 50: self.pending_messages = self.pending_messages[-50:] def _process_queue(self): if not self.send_paused and not self.is_sending and self.pending_messages: # 继续队列处理逻辑 pass
3. 实现客户端ACK确认机制
让客户端收到每条消息后返回确认信号,后端仅在收到ACK后才发送下一条消息,确保消息确实被接收,彻底避免无意义的缓冲区堆积。
示例流程:
- 后端发送消息附带唯一标识:
{"id": 123, "data": "xxx"} - 客户端收到后发送确认:
{"ack": 123} - 后端收到ACK后,再发送队列中的下一条消息
4. 慢连接检测与降级
定期统计消息发送吞吐量,当发现连续周期内发送速率远低于生产速率时,采取以下措施:
- 降低消息生产频率(业务允许的情况下)
- 主动断开连接并告知客户端重连
- 合并多条消息为批量发送,减少IO开销
内容的提问来源于stack exchange,提问作者themozel
相关产品推荐
相关产品推荐

