小Cの已经记不起来的博客

WebSocket 实时通信与 熔断 应用

最近在做一个实时行情推送的项目,前端用 WebSocket 连后端,后端从 Kafka 拿数据再推给前端。本来觉得挺简单的,结果上线第一天就翻车了。某个后端实例因为内存泄漏卡住了,前端一看连接断了,就开始疯狂重连,一秒一次,网关连接数直接飙到几万,最后整个服务都挂了。当时就想着,这玩意儿得加个熔断啊。

WebSocket 的长连接没那么简单

WebSocket 跟普通 HTTP 不一样,它是长连接,建立之后双方可以一直互相发消息。听起来很美,但维护长连接的成本其实挺高的。网络抖动、服务端重启、负载均衡超时,任何一个环节出问题,连接就断了。客户端一般都会写个自动重连,比如:

const ws = new WebSocket('wss://example.com/ws');
ws.onclose = () => {
  setTimeout(() => connect(), 1000);
};

这段代码在小规模下没问题,但如果服务端已经扛不住了,所有客户端都在一秒后同时重连,那就是雪崩。重连风暴比原来的故障还可怕。

熔断到底在断什么

熔断器(Circuit Breaker)这个概念在微服务里很常见,通常有三种状态:闭合(CLOSED)、打开(OPEN)、半开(HALF_OPEN)。闭合时正常放行请求,失败次数达到阈值就打开,打开后直接拒绝请求,不再往故障服务上撞。等过了一段时间,进入半开状态,放一个请求过去试探,成功就闭合,失败就继续打开。

在 WebSocket 场景下,熔断可以放在两个地方:客户端和服务端代理层。客户端熔断主要是控制重连频率,服务端熔断主要是限制连接数和后端调用。

客户端重连的熔断改造

我一开始用的是 Python 的 websockets 库,重连逻辑写得跟上面 JS 差不多。后来加了一个简单的熔断器:

import asyncio
import websockets
from datetime import datetime, timedelta

class CircuitBreaker:
    def __init__(self, failure_threshold=5, recovery_timeout=30):
        self.failure_threshold = failure_threshold
        self.recovery_timeout = recovery_timeout
        self.failures = 0
        self.state = "CLOSED"
        self.last_failure_time = None

    def record_failure(self):
        self.failures += 1
        if self.failures >= self.failure_threshold:
            self.state = "OPEN"
            self.last_failure_time = datetime.now()
            print("熔断器打开,暂停重连")

    def record_success(self):
        self.failures = 0
        self.state = "CLOSED"

    def allow_request(self):
        if self.state == "CLOSED":
            return True
        if self.state == "OPEN":
            if datetime.now() - self.last_failure_time > timedelta(seconds=self.recovery_timeout):
                self.state = "HALF_OPEN"
                return True
            return False
        if self.state == "HALF_OPEN":
            return True

然后连接循环变成这样:

async def connect_with_breaker():
    breaker = CircuitBreaker()
    while True:
        if not breaker.allow_request():
            await asyncio.sleep(1)
            continue
        try:
            async with websockets.connect("ws://localhost:8765") as ws:
                breaker.record_success()
                await handle_messages(ws)
        except Exception as e:
            print(f"连接失败: {e}")
            breaker.record_failure()
            await asyncio.sleep(1)

这样连续失败 5 次后,熔断器打开,客户端会安静 30 秒再尝试。半开状态只放一个连接进去,避免了重连风暴。当然实际用的时候还得加个指数退避,比如 sleep(min(2 ** failures, 60)),不然半开的时候还是可能撞。

服务端代理层的熔断配置

客户端改了还不够,服务端也得限制。我们用的是 Envoy 做网关,它自带的 circuit_breakers 对 WebSocket 很管用:

circuit_breakers:
  thresholds:
    - priority: DEFAULT
      max_connections: 1000
      max_pending_requests: 100
      max_requests: 1000
      max_retries: 3

max_connections 直接限制单个上游集群的最大连接数,超过就直接拒绝,不会让连接堆积。如果你用 Nginx,虽然没有原生的熔断,但可以用 limit_connmax_conns 凑合:

upstream backend {
    server 127.0.0.1:8080 max_conns=500;
    server 127.0.0.1:8081 max_conns=500;
}

limit_conn_zone $binary_remote_addr zone=perip:10m;
limit_conn perip 10;

不过 Nginx 的 max_conns 只对上游生效,而且配置重载时不太灵活。真要搞熔断,还是 Envoy 或者 Istio 更省心。

几个容易踩的坑

第一个是心跳不能太频繁。我一开始设置 5 秒没收到消息就判定断开,结果网络稍微抖一下熔断器就开了。后来改成 30 秒,世界清净了。

第二个是半开状态一定要加锁。多个协程同时判断 allow_request(),可能会同时进入半开,然后一起失败,熔断器又打开了。简单用个 asyncio.Lock 就行。

第三个是熔断恢复时间要配合业务。比如行情推送,30 秒不推数据用户能接受,但如果是聊天室,可能 10 秒就炸了。这个得自己调。

第四个是别只做客户端熔断。服务端该限流还得限流,该扩容还得扩容。熔断只是防止雪崩,不是银弹。

最后

WebSocket 实时通信加上熔断之后,至少不会因为一个实例挂掉就把整个集群带走。代码其实不复杂,关键是理解熔断的状态机和重连的节奏。如果你也在搞实时推送,不妨试试。不出问题的话就没有问题了,出了问题再回来调阈值。

评论

还没有评论。

发表评论

提交后评论将经过自动审核,审核通过后公开展示。

未在播放