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
相关产品推荐
相关产品推荐

