Python如何测量等待ws.recv()处理的WebSocket消息队列长度
WebSocket消息队列堆积测量与最大数据流测算方案
一、队列等待消息数测量实现
你使用的是websocket-client库的同步客户端,该库没有直接暴露内核缓冲区/内部消息队列长度的公开API,可通过非阻塞批量拉取的方式统计堆积数量,对原有逻辑的性能影响可控制在微秒级,完全适配你的低延迟场景需求:
- 核心思路:每次正常拉取1条消息处理完成后,用非阻塞模式尝试拉取所有已经到达内核缓冲区的消息,拉取到的总条数就是当前队列的堆积数,统计完成后再按顺序处理这些消息即可。
二、最大可消费数据流测算方法
按以下步骤执行压测即可得到准确上限:
- 逐步增加订阅的行情数据流数量,每轮压测持续10分钟以上覆盖消息突发场景
- 代码中埋点记录两个指标:队列堆积峰值条数、单条消息平均处理耗时
- 当出现
堆积峰值条数 * 单条处理耗时 > 10ms时,当前的数据流数量就是你要的上限值
三、代码修改示例
仅需要修改_listen方法即可实现堆积统计,修改后代码如下:
import select import time import socket # 其他原有代码不变 def _listen(self): self.keepalive.start() # 统计用变量,可根据需要上报到监控 self.max_queue_size = 0 self.avg_process_time = 0 self.msg_count = 0 # 可选:扩大接收缓冲区避免突发流量丢包,这里设为4MB self.ws.sock.setsockopt(socket.SOL_SOCKET, socket.SO_RCVBUF, 1024*1024*4) while not self.stop: try: # 先阻塞拉取第一条消息,保证没有消息时不占CPU data = self.ws.recv() start_process = time.perf_counter() msg = json.loads(data) self.on_message(msg) # 统计单条处理耗时 process_cost = time.perf_counter() - start_process self.msg_count +=1 self.avg_process_time = (self.avg_process_time * (self.msg_count-1) + process_cost) / self.msg_count # 开始统计当前队列堆积数 queue_size = 0 pending_msgs = [] # 用select判断是否有可读数据,超时设为0直接返回 while True: read_list, _, _ = select.select([self.ws.sock], [], [], 0) if not read_list: break # 有数据就非阻塞拉取 self.ws.settimeout(0) try: pending_data = self.ws.recv() pending_msgs.append(json.loads(pending_data)) queue_size +=1 except: # 没有更多数据就退出 break finally: # 恢复阻塞模式 self.ws.settimeout(None) # 更新堆积峰值 if queue_size > self.max_queue_size: self.max_queue_size = queue_size # 处理堆积的消息 for pending_msg in pending_msgs: start_process = time.perf_counter() self.on_message(pending_msg) process_cost = time.perf_counter() - start_process self.msg_count +=1 self.avg_process_time = (self.avg_process_time * (self.msg_count-1) + process_cost) / self.msg_count except ValueError as e: self.on_error(e) except Exception as e: self.on_error(e)
注意事项
- 统计逻辑可按需要新增日志打印或监控上报,不需要时可直接注释掉,不影响原有业务逻辑运行
- 若压测过程中出现丢包情况,可进一步调大socket接收缓冲区的配置值
内容的提问来源于stack exchange,提问作者Chris
相关产品推荐
相关产品推荐

