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_conn 和 max_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 实时通信加上熔断之后,至少不会因为一个实例挂掉就把整个集群带走。代码其实不复杂,关键是理解熔断的状态机和重连的节奏。如果你也在搞实时推送,不妨试试。不出问题的话就没有问题了,出了问题再回来调阈值。
评论
还没有评论。
发表评论
提交后评论将经过自动审核,审核通过后公开展示。