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

Python websocket库如何实现多消息监听器 对标JS的add_event_listener功能

核心实现思路

你可以自己实现一套轻量的订阅分发中间层替代原生单一的on_message回调,不需要依赖全局变量,实现多监听者的注册、规则匹配、自动订阅/取消订阅能力,完全可以对标JS的add_event_listener使用体验。

完整实现代码
import websocket
import json
import threading
from uuid import uuid4

class MultiListenerWebSocketClient:
    def __init__(self, socket_url):
        self.socket_url = socket_url
        self.ws = None
        # 存储所有订阅项,key为订阅ID,value存储过滤规则、处理函数、取消订阅请求
        self.subscriptions = {}
        self.lock = threading.Lock()

    def _on_message(self, ws, message):
        # 可根据业务调整是否需要JSON解析,纯文本消息可跳过该步
        try:
            msg_data = json.loads(message)
        except json.JSONDecodeError:
            msg_data = message
        
        # 遍历所有订阅项,匹配规则则触发对应处理函数
        with self.lock:
            for sub_info in self.subscriptions.values():
                if sub_info["filter"](msg_data):
                    # 处理逻辑如果耗时较长,可改成丢到线程池执行避免阻塞消息接收
                    sub_info["handler"](msg_data)

    def _on_open(self, ws):
        # 连接建立的自定义逻辑可在这里扩展
        pass

    def _on_close(self, ws, close_status_code, close_msg):
        # 连接关闭的自定义逻辑可在这里扩展
        pass

    def add_message_listener(self, filter_func, handler, subscribe_request=None):
        """
        添加消息监听器
        :param filter_func: 过滤函数,入参为收到的消息,返回True则将消息传给当前handler处理
        :param handler: 消息处理函数,入参为匹配成功的消息
        :param subscribe_request: 可选,订阅请求,添加监听时自动发送给服务端,可携带取消订阅的配置
        :return: 订阅ID,用于后续取消订阅
        """
        sub_id = str(uuid4())
        with self.lock:
            self.subscriptions[sub_id] = {
                "filter": filter_func,
                "handler": handler,
                "unsubscribe_request": subscribe_request.pop("unsubscribe", None) if subscribe_request else None
            }
        # 连接已建立的场景下自动发送订阅请求
        if subscribe_request and self.ws and self.ws.sock and self.ws.sock.connected:
            self.ws.send(json.dumps(subscribe_request))
        return sub_id

    def remove_message_listener(self, sub_id):
        """
        取消消息监听
        :param sub_id: 添加监听时返回的订阅ID
        """
        with self.lock:
            sub_info = self.subscriptions.pop(sub_id, None)
        # 自动发送取消订阅请求
        if sub_info and sub_info["unsubscribe_request"] and self.ws and self.ws.sock and self.ws.sock.connected:
            self.ws.send(json.dumps(sub_info["unsubscribe_request"]))

    def start(self):
        """启动websocket连接,后台线程运行不阻塞主线程操作"""
        self.ws = websocket.WebSocketApp(
            self.socket_url,
            on_open=self._on_open,
            on_close=self._on_close,
            on_message=self._on_message
        )
        threading.Thread(target=self.ws.run_forever, daemon=True).start()
使用示例
# 初始化客户端并启动连接
SOCKET = "wss://你的websocket服务地址"
client = MultiListenerWebSocketClient(SOCKET)
client.start()

# 定义监听器1:仅处理param1对应的消息
def param1_handler(msg):
    print("收到param1的消息:", msg)
# 自定义过滤规则
param1_filter = lambda msg: isinstance(msg, dict) and msg.get("param") == "param1"
# 配置订阅、取消订阅请求
param1_sub_req = {
    "method": "SUBSCRIBE","params":["param1"],"id": 1,
    "unsubscribe": {"method": "UNSUBSCRIBE","params":["param1"],"id": 1}
}
# 添加监听,拿到订阅ID
sub_id1 = client.add_message_listener(param1_filter, param1_handler, param1_sub_req)

# 定义监听器2:仅处理param2对应的消息
def param2_handler(msg):
    print("收到param2的消息:", msg)
param2_filter = lambda msg: isinstance(msg, dict) and msg.get("param") == "param2"
param2_sub_req = {
    "method": "SUBSCRIBE","params":["param2"],"id": 2,
    "unsubscribe": {"method": "UNSUBSCRIBE","params":["param2"],"id": 2}
}
sub_id2 = client.add_message_listener(param2_filter, param2_handler, param2_sub_req)

# 不需要param1的消息时取消订阅即可
# client.remove_message_listener(sub_id1)
方案优势
  • 完全解耦各监听逻辑,无需全局变量存储消息
  • 每个监听者可自定义任意过滤规则,不限于匹配订阅参数,可灵活扩展
  • 自动处理订阅、取消订阅的请求发送,业务侧无需单独调用发送逻辑
  • 加线程锁保证多线程场景下订阅列表操作的安全性
  • 支持任意数量的监听者同时工作

内容的提问来源于stack exchange,提问作者Abdullahi Ibrahim

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.01 10:27:05