You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

Python如何测量等待ws.recv()处理的WebSocket消息队列长度

WebSocket消息队列堆积测量与最大数据流测算方案

一、队列等待消息数测量实现

你使用的是websocket-client库的同步客户端,该库没有直接暴露内核缓冲区/内部消息队列长度的公开API,可通过非阻塞批量拉取的方式统计堆积数量,对原有逻辑的性能影响可控制在微秒级,完全适配你的低延迟场景需求:

  • 核心思路:每次正常拉取1条消息处理完成后,用非阻塞模式尝试拉取所有已经到达内核缓冲区的消息,拉取到的总条数就是当前队列的堆积数,统计完成后再按顺序处理这些消息即可。

二、最大可消费数据流测算方法

按以下步骤执行压测即可得到准确上限:

  1. 逐步增加订阅的行情数据流数量,每轮压测持续10分钟以上覆盖消息突发场景
  2. 代码中埋点记录两个指标:队列堆积峰值条数、单条消息平均处理耗时
  3. 当出现 堆积峰值条数 * 单条处理耗时 > 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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.09.25 23:45:01