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

多客户端WebSocket流端点(Python):并发阻塞问题及优化问询

解决Sanic WebSocket多客户端阻塞问题,以及大厂流式服务实现思路

你的问题根源在于单个WebSocket连接的处理协程占用了事件循环的全部时间——虽然用了await,但你的while True循环里没有任何让出执行权的操作(比如延迟、等待外部事件),导致事件循环无法调度其他客户端的连接请求。下面分两部分给你拆解解决方案和大厂的实现思路:

一、快速修复你的Sanic代码,支持多客户端并发

我们可以通过维护活跃连接池+独立广播任务的方式,让单个事件循环就能处理大量并发WebSocket连接,完全不需要增加多个worker。

修改后的代码如下:

from sanic import Sanic
import asyncio

app = Sanic("WebSocketBroadcast")
ws_routes = app

# 用集合存储所有活跃的WebSocket连接(Sanic事件循环单线程,集合操作安全)
active_connections = set()

@ws_routes.websocket("/hello")
async def handle_websocket(request, ws):
    # 新连接加入集合
    active_connections.add(ws)
    try:
        # 保持连接存活,同时监听客户端的关闭信号
        while True:
            # 接收客户端消息(如果不需要处理输入,可改为await asyncio.sleep(3600),但监听关闭更合理)
            msg = await ws.receive()
            if msg.type == "websocket.close":
                break
    finally:
        # 连接关闭时从集合移除
        active_connections.discard(ws)

# 后台独立任务:定时给所有活跃连接广播消息
async def broadcast_loop():
    while True:
        # 迭代集合副本,避免迭代过程中集合被修改(比如连接关闭)
        for ws in list(active_connections):
            try:
                await ws.send("hello")
            except Exception as e:
                # 发送失败(比如连接已断),自动清理无效连接
                print(f"发送消息失败,移除无效连接: {str(e)}")
                active_connections.discard(ws)
        # 每隔1秒发送一次,强制让出执行权,让事件循环处理其他连接
        await asyncio.sleep(1)

if __name__ == "__main__":
    # 注册后台广播任务
    app.add_task(broadcast_loop())
    app.run(host="0.0.0.0", port=8000)

代码说明:

  1. 连接池管理:用active_connections集合跟踪所有在线的WebSocket连接,连接建立时加入,关闭时自动移除。
  2. 独立广播任务:把消息发送逻辑从单个连接的处理函数中抽离,由单独的协程负责批量发送,每次发送后用asyncio.sleep(1)让出事件循环,确保其他连接的请求能被及时处理。
  3. 容错处理:发送消息时捕获异常,自动清理无效连接,避免集合中残留死连接占用资源。

这样修改后,单worker就能轻松处理成百上千个并发客户端,完全不需要为每个客户端单独开worker。

二、Binance、Twitter这类大厂的流式WebSocket服务实现思路

这类平台能支持百万级并发连接,核心是异步IO事件驱动+分布式架构,具体细节包括:

1. 基于轻量级协程的连接模型

他们用异步IO框架(比如Node.js、Go的goroutine、Python的asyncio),每个WebSocket连接对应一个轻量级协程(而非进程/线程)。单个进程就能处理数万甚至十万级别的并发连接——协程的内存开销只有几KB,远低于线程/进程的资源消耗。

2. 中心化的连接管理与消息广播

  • 维护全局连接池(集群场景下会用分布式存储/消息队列同步连接状态),当有新的实时数据(比如行情、推文)时,通过广播机制推送给所有订阅的客户端,而非让客户端主动拉取。
  • 消息推送会做批量优化,比如合并相同的消息,减少IO传输开销。

3. 水平扩展与负载均衡

  • 前端部署负载均衡器(比如Nginx、HAProxy),将客户端连接均匀分发到后端多台服务器节点。
  • 用消息队列(比如Kafka、Redis Pub/Sub)实现跨节点的消息同步:当某台节点收到数据源的新数据,会把数据发送到消息队列,其他节点从队列中获取数据后推送给自己的客户端,保证全集群的消息一致。

4. 心跳与连接健康检测

  • 客户端和服务器之间定期发送心跳包,检测连接状态,自动清理死连接,释放资源。
  • 客户端实现自动重连逻辑,确保网络波动后能快速恢复连接。

5. 性能优化

  • 用二进制协议(比如Protobuf)替代JSON,减少消息体积,提升传输效率。
  • 对高频消息做限流或合并处理,避免服务器和客户端被过多消息压垮。

总结

你之前的问题是因为单个连接的循环没有让出事件循环,导致其他连接无法被调度。用连接池+独立广播任务的方式就能解决单worker的多客户端并发问题。如果未来需要支持更大规模的用户,再逐步引入负载均衡、消息队列等分布式组件,就能达到大厂的服务水平。

内容的提问来源于stack exchange,提问作者Halcyon Abraham Ramirez

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 08:47:15