Python实现SocketIO与WebSocket服务器消息转发桥接的问题
解决SocketIO与WebSocket桥接的同步异步交互及连接复用问题
核心问题拆解
你遇到的本质问题是同步事件回调(SocketIO的on_message)无法直接调用异步WebSocket操作,同时需要复用单个WebSocket连接避免频繁建立连接的开销。下面是针对性的解决方案,附上完整可运行代码及关键逻辑说明。
修复后完整代码
import socketio import asyncio import websockets import json from typing import Dict, Any # 配置服务器地址 SOCKETIO_ADDR = 'http://localhost:1234' WEBSOCKET_ADDR = 'ws://localhost:1235' # 使用异步SocketIO客户端,支持异步回调 sio = socketio.AsyncClient() # 消息队列:作为SocketIO回调与WebSocket任务的中转通道 message_queue = asyncio.Queue() # SocketIO消息回调(异步版本,可直接操作队列) @sio.on('message') async def on_socketio_message(data: Dict[str, Any]): message_id = data.get('messageId') value = data.get('value') if message_id == "requiredTag": event = {"sender": "myclients_name", message_id: value} # 将待发送消息放入队列,交由WebSocket任务处理 await message_queue.put(json.dumps(event)) async def websocket_bridge(): # 建立单个WebSocket连接并复用 async with websockets.connect(WEBSOCKET_ADDR) as websocket: # 启动两个并行子任务:接收WebSocket消息、发送队列中的消息 receive_task = asyncio.create_task(receive_from_websocket(websocket)) send_task = asyncio.create_task(send_to_websocket(websocket)) # 等待任一任务终止(通常是连接断开),清理剩余任务 done, pending = await asyncio.wait( [receive_task, send_task], return_when=asyncio.FIRST_COMPLETED ) for task in pending: task.cancel() async def receive_from_websocket(websocket): # 持续接收WebSocket消息并转发到SocketIO while True: response = await websocket.recv() response_dict = json.loads(response) response_list = extract_values(response_dict, [], []) for key, value in response_list: await sio.emit('message_from_websocketserver', {'messageId': key, 'value': value}) async def send_to_websocket(websocket): # 持续从队列取消息并通过WebSocket发送 while True: message = await message_queue.get() await websocket.send(message) message_queue.task_done() # 补全你提到的extract_values解析函数(如果未实现) def extract_values(data, path, result): if isinstance(data, dict): for k, v in data.items(): extract_values(v, path + [k], result) elif isinstance(data, list): for idx, item in enumerate(data): extract_values(item, path + [str(idx)], result) else: result.append(('.'.join(path), data)) return result async def main(): # 先连接SocketIO服务器 await sio.connect(SOCKETIO_ADDR) # 启动WebSocket桥接逻辑 await websocket_bridge() # 断开SocketIO连接 await sio.disconnect() if __name__ == "__main__": asyncio.run(main())
关键解决逻辑
改用异步SocketIO客户端
使用socketio.AsyncClient替代同步客户端,让on_message支持异步回调,直接与asyncio队列交互,避免同步/异步代码冲突。用asyncio.Queue做消息中转
SocketIO回调将待发送消息放入队列,WebSocket的发送任务从队列取消息发送,实现两类事件的解耦,同时保证复用同一个WebSocket连接。单WebSocket连接复用
在websocket_bridge中建立单个连接,同时启动接收、发送两个子任务,同一个连接既处理上游消息接收,又处理下游消息发送,彻底避免频繁建连的开销。任务生命周期管理
通过asyncio.create_task管理并行任务,当连接断开时自动取消未完成的任务,确保资源正常释放。
兼容同步SocketIO客户端的方案
如果必须使用同步SocketIO客户端,可通过run_coroutine_threadsafe在同步回调中提交异步任务:
# 同步SocketIO客户端的回调函数 @sio.on('message') def on_socketio_message(data): message_id = data.get('messageId') value = data.get('value') if message_id == "requiredTag": event = {"sender": "myclients_name", message_id: value} # 在同步函数中向事件循环提交异步任务 loop = asyncio.get_event_loop() loop.call_soon_threadsafe(message_queue.put_nowait, json.dumps(event)) # 主函数调整为多线程运行 def main(): loop = asyncio.new_event_loop() asyncio.set_event_loop(loop) # 启动SocketIO客户端在后台线程 import threading threading.Thread(target=sio.connect, args=(SOCKETIO_ADDR,), daemon=True).start() # 运行WebSocket桥接逻辑 loop.run_until_complete(websocket_bridge())
内容的提问来源于stack exchange,提问作者apfreelance
相关产品推荐
相关产品推荐

