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

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())

关键解决逻辑

  1. 改用异步SocketIO客户端
    使用socketio.AsyncClient替代同步客户端,让on_message支持异步回调,直接与asyncio队列交互,避免同步/异步代码冲突。

  2. 用asyncio.Queue做消息中转
    SocketIO回调将待发送消息放入队列,WebSocket的发送任务从队列取消息发送,实现两类事件的解耦,同时保证复用同一个WebSocket连接。

  3. 单WebSocket连接复用
    在websocket_bridge中建立单个连接,同时启动接收、发送两个子任务,同一个连接既处理上游消息接收,又处理下游消息发送,彻底避免频繁建连的开销。

  4. 任务生命周期管理
    通过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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.27 08:13:10