当我们聊到 WebSocket 服务的时候,很多人第一反应就是“能推消息、双向通信、比轮询好”。但如果这个服务要面对成千上万的连接,问题就不那么简单了。连接一多,内存就像漏水的地桶,CPU 也跟着飙高,一不小心进程就挂了。今天咱们不拽术语,就用大白话,拿 Actix Web 这个 Rust 框架,聊聊怎么把这种服务做得稳一点,重点说清楚“背压”和“内存管理”这两件事。
一、先聊聊 WebSocket 的“胃口”有多大
WebSocket 和普通 HTTP 不一样,它一旦建立连接,就是一条长期开着的通道。客户端不断发消息,服务端也不断推数据。每条消息虽然不大,但架不住连接基数大。比如你有 1 万个用户,每个人每秒钟发一条 1KB 的消息,那每秒就要处理 10MB 数据。这还不算头信息、缓冲区和连接状态的开销。
更麻烦的是,WebSocket 消息是异步到达的。服务端收到一条消息,马上就得处理,处理不完就得排队。队列放在内存里,如果来不及消费,队列就越来越长。内存一点点被吃掉,最后就像堵车一样,整个服务动弹不得。Rust 虽然内存安全,但安全不代表不占内存,也不代表不会撑爆。所以,我们得主动控制“背压”,让那些来不及处理的消息不要一股脑涌进来。
二、Actix Web 里 WebSocket 的基本玩法
先说一个最简单的例子,让你知道 Actix Web 里怎么接 WebSocket。假设我们做一个聊天服务,客户端连上来就能互相发消息。当然,真正的聊天要搞房间和转发逻辑,这里只是为了演示基本结构。
技术栈:Rust + Actix Web + actix-ws
// 引入必要的依赖。记住,Cargo.toml 里要有 actix-web 和 actix-ws
use actix_web::{web, App, HttpServer, HttpResponse};
use actix_web::get;
use actix_ws::Message;
#[get("/ws")]
async fn ws_handler(req: actix_web::HttpRequest, stream: actix_web::web::Payload) -> actix_web::Result<HttpResponse> {
// 这是 Actix Web 处理 WebSocket 升级的标准调用
let (response, mut session, mut msg_stream) = actix_ws::handle(&req, stream)?;
// 先启动一个异步任务,专门处理来自客户端的消息
actix_web::rt::spawn(async move {
// 循环读取客户端发来的消息
while let Some(Ok(msg)) = msg_stream.next().await {
// 如果是文本消息,就把内容回显给客户端
match msg {
Message::Text(text) => {
let _ = session.text(text).await; // 注意这里 .await,可能会被阻塞
}
Message::Close(reason) => {
println!("连接关闭: {:?}", reason);
break;
}
_ => { /* 忽略二进制及其他类型 */ }
}
}
});
// 立刻返回握手响应,让连接升级
Ok(response)
}
#[actix_web::main]
async fn main() -> std::io::Result<()> {
HttpServer::new(|| {
App::new()
.service(ws_handler)
})
.bind("127.0.0.1:8080")?
.run()
.await
}
这个例子看着没问题吧?其实里面藏着一个隐患:session.text(text).await 等待发送完成,而发送是会阻塞的。如果客户端接收慢,服务端发一条消息要等很久。一个连接等没事,一万个连接都在等,系统就卡死了。所以,我们不能这样简单地把发送和接收放在同一个循环里。
三、什么是背压?为什么它像“水位线”
背压这个词听起来吓人,其实就是“上游的推力太大,下游跟不上,得往回顶一顶”。想象一根水管,水龙头开得很大,但出水口太细,管子会爆。背压控制就是装一个阀门,当水管里压力太大时,自动关小水龙头。
在 WebSocket 服务里,背压体现在两个方向:
- 客户端发消息太快,服务端处理不过来。这时不能让客户端无限地发,要么拒绝,要么让客户端等一等。
- 服务端推消息太快,客户端接收不过来。比如服务端批量推送数据,如果客户端带宽小,服务端就会堆积。
Actix Web 提供了一些机制来感知这种压力。比如 session.text() 返回的 future 会等待消息真正写入传输层。如果传输层的 buffer 满了,这个 future 就会一直挂起。如果我们不加限制,挂起的任务会积累,导致内存爆炸。
怎么办呢?最常见的办法是给每个连接设置一个“待发送队列”的上限。超过上限,就丢掉消息,或者关闭连接。这样,压力就不会无限积累。
四、在 Actix 里做背压控制:代码实战
4.1 使用有界通道来收消息
Actix Web 的 MessageStream 本身可以接在通道上。我们可以用 tokio::sync::mpsc 创建一个有界的通道,把收到的消息塞进去,然后由另一个任务慢慢处理。通道的数量限制就是背压的阈值。
技术栈:Rust + Actix Web + actix-ws + tokio
use actix_web::{web, App, HttpServer, HttpResponse};
use actix_web::get;
use actix_ws::Message;
use tokio::sync::mpsc;
#[get("/ws")]
async fn ws_handler(req: actix_web::HttpRequest, stream: actix_web::web::Payload) -> actix_web::Result<HttpResponse> {
let (response, mut session, mut msg_stream) = actix_ws::handle(&req, stream)?;
// 创建一个有界通道,最多缓存 128 条消息。这就是我们的“背压阀门”
let (tx, mut rx) = mpsc::channel::<Message>(128);
// 任务1:把客户端发来的每条消息,通过 tx 发到通道里
actix_web::rt::spawn(async move {
while let Some(Ok(msg)) = msg_stream.next().await {
// 如果发送到通道失败(比如接受端被关闭了),直接退出发送循环
if tx.send(msg).await.is_err() {
break;
}
}
});
// 任务2:从通道的另一端取消息,做真正的业务处理
actix_web::rt::spawn(async move {
while let Some(msg) = rx.recv().await {
match msg {
Message::Text(text) => {
// 这里模拟业务处理,比如计算、存储、转发
let reply = format!("收到: {}", text);
// 发送回复。如果客户端太慢,.await 会阻塞当前任务
// 但因为每个连接都只占一个任务,阻塞不会影响到其他连接
let _ = session.text(reply).await;
}
Message::Close(_) => {
println!("对方要求关闭");
break;
}
_ => {}
}
}
});
Ok(response)
}
有界通道的好处是,如果客户端疯狂发送消息,而业务处理很慢,通道里最多存 128 条。再往后的消息,tx.send().await 会一直等,但注意:发送端任务也会被阻塞。实际上,如果客户端一直发,服务端的任务会挂起,不再从 socket 读取新消息,这本身就形成了一种“反推力”,让 TCP 层的缓冲区慢慢变满,客户端就会被减速。这就是背压的本质。
4.2 限制单条消息大小
除了限制数量,还要限制单条消息的体积。有些人会发一个几十兆的文本,直接把你内存干爆。Actix Web 中,我们可以通过 actix_ws::MessageStream 配置消息大小限制吗?实际上,actix_ws 依赖底层的 futures 流,我们可以自己判断消息大小。
技术栈:Rust + Actix Web + actix-ws
// 统计当前消息的字节数,如果超过阈值就拒绝
const MAX_MESSAGE_BYTES: usize = 64 * 1024; // 64KB
async fn handle_message(msg: Message, session: &mut actix_ws::Session) -> bool {
match msg {
Message::Text(text) => {
// text 是 Bytes 类型,拿到它的长度
if text.len() > MAX_MESSAGE_BYTES {
println!("消息太大,拒绝处理");
let _ = session.close(None).await; // 直接关闭连接
return false;
}
let reply = format!("长度正常: {}", text.len());
let _ = session.text(reply).await;
true
}
Message::Close(_) => {
println!("客户端关闭");
false
}
_ => true,
}
}
然后在接收循环里调用即可。这样,即便有人恶意发超大消息,我们也不会把大对象塞进通道,而是直接断开,保护了内存。
4.3 定时量内存和连接数
除了背压,我们还得知道自己的服务到底有多少连接,占了多少内存。Actix Web 提供了一些数据?其实没有内置统计,但我们可以自己维护一个连接计数器。
技术栈:Rust + Actix Web + std::sync::Arc
use std::sync::atomic::{AtomicUsize, Ordering};
use std::sync::Arc;
// 全局连接计数器
struct AppState {
conn_count: AtomicUsize,
}
async fn ws_handler(
req: actix_web::HttpRequest,
stream: actix_web::web::Payload,
state: web::Data<Arc<AppState>>,
) -> actix_web::Result<HttpResponse> {
let conn_count = state.conn_count.fetch_add(1, Ordering::SeqCst);
println!("当前连接数: {}", conn_count + 1);
// 如果连接数超过预设最大值,直接拒绝升级
const MAX_CONN: usize = 100_000;
if conn_count >= MAX_CONN {
// 先递减恢复,然后返回 503
state.conn_count.fetch_sub(1, Ordering::SeqCst);
return Ok(HttpResponse::ServiceUnavailable().finish());
}
let (response, mut session, mut msg_stream) = actix_ws::handle(&req, stream)?;
actix_web::rt::spawn(async move {
// 在连接结束时,递减计数器
let state_ref = state.clone();
// 用 Guard 模式,确保在连接退出时自动减一
struct Guard(Arc<AppState>);
impl Drop for Guard {
fn drop(&mut self) {
self.0.conn_count.fetch_sub(1, Ordering::SeqCst);
}
}
let _guard = Guard(state_ref.clone());
// 正常处理消息
while let Some(Ok(msg)) = msg_stream.next().await {
// 处理逻辑略
let _ = session.text("pong").await;
}
});
Ok(response)
}
这样我们就能实时掌握连接数,并且设置一个硬阈值,防止程序被大量连接拖垮。内存方面,Rust 没有内置的“查看当前内存”API,但我们可以通过系统监控工具来观察。思路是:不要让每一个连接都存着巨大的缓冲,尽量做到“来一条消息,处理一条消息,丢弃一条消息”。
五、内存管理的几个朴素习惯
5.1 避免无界提交
前面说的通道是有界的,这是最重要的习惯。很多人喜欢用 VecDeque 或者 Vec 来缓存消息,那是在给自己挖坑。无界集合在压力下会无限增长,直到内存耗尽。Rust 的标准库没有提供有界集合,但我们可以用 tokio::sync::mpsc 或者 crossbeam_channel::bounded。
5.2 共享状态要尽量小
如果多个连接需要共享数据,比如房间列表、在线用户表,可以考虑用 Arc<Mutex<HashMap>> 或者 dashmap。但是要注意,锁不住的时候会把线程挂起,等待锁的任务也可能堆积。一个更稳妥的方式是用 actix_web::web::Data + 一个无锁的数据结构,或者把数据放到 Redis 这类外部服务里。
5.3 合理设置 TCP 缓冲
操作系统会给每个 TCP 连接分配发送缓冲和接收缓冲。你可以通过 socket2 来改这些值,但更简单的是在系统层面调优。在 Actix Web 里,我们可以通过 HttpServer::max_connection_rate 限制每秒新连接数,通过 max_connections 限制总连接数。虽然这不是内存管理本身,但能变相降低内存压力。
5.4 及时清理死连接
WebSocket 连接不一定会有正常的关闭帧。如果客户端断电、断网,服务端可能很久都感知不到。所以我们要启用心跳检测。Actix Web 没有内置的心跳,但我们可以用 actix_ws::Session::ping 定期发送 ping 帧,如果几次没收到 pong,就关闭连接。这是一种主动的内存清理,把占着位置不工作的连接踢掉。
技术栈:Rust + Actix Web + actix-ws + tokio time
use actix_ws::Message;
use tokio::time::{interval, Duration};
async fn heartbeat(session: &mut actix_ws::Session) {
let mut ticker = interval(Duration::from_secs(10));
loop {
ticker.tick().await;
// 发送 ping,如果发送失败说明连接已经断开
if session.ping(b"").await.is_err() {
println!("心跳失败,关闭连接");
let _ = session.close(None).await;
break;
}
}
}
// 在使用时,你需要在处理消息的循环中同时select心跳任务的完成
// 这里简略展示,实际可以用 tokio::select!
六、应用场景与优缺点
6.1 适合的场景
高并发 WebSocket 服务非常适合实时消息推送、在线聊天、股票行情、多人小游戏、协同编辑。这些场景都有一个特点:连接多、消息频繁,但每条消息不大。如果用在物联网设备状态上报,也很有价值,因为设备往往是大量低功率的终端,连接稳定性不如浏览器,更需要心跳和背压保护。
6.2 使用 Actix Web 的优点
Actix Web 基于 actor 模型,但说起来更像是一个异步运行时。它的性能在 Rust 生态里数一数二,内存占用相对可控。它提供了 actix-ws 这个专门的 WebSocket 库,处理协议细节很省心。Rust 本身没有 GC,你不需要担心触发全局垃圾回收导致卡顿,这对于高并发是巨大的优势。
6.3 需要注意的缺点
Actix Web 的学习曲线有点陡。如果你之前习惯 Node.js 或者 Go,会觉得 Rust 的语法很繁琐,尤其是生命周期和所有权。actix-ws 的 API 还在变化,不同版本之间兼容性一般。另外,Rust 的 async 并发模型里,Send 约束有时候会折磨人,想把一个消息发送器克隆到多个任务里,得花一些时间才能搞定。
七、注意事项
- 不要依赖默认的配置。以为 Actix Web 很牛就什么都不管,连接数、消息大小、队列长度不设限制,早晚出事。
- 不要在锁里做耗时操作。比如持有一个
Mutex去写数据库,会把其他连接统统堵住。尽量用锁只保护临界区,或者把写入操作放到异步任务里。 - 正确处理
close帧。客户端发来关闭信号时,服务端要把连接清理干净,别再往里发消息了。 - 使用压缩功能要慎用。WebSocket 的 permessage-deflate 扩展能省带宽,但会消耗 CPU 和内存,在高并发时可能得不偿失。如果不需要,最好不要开启。
- 测试压力不能少。上线前用
websocat或者oha模拟几千个并发连接,观察内存和 CPU。别在乎理论峰值,实际跑一下最靠谱。
八、总结
高并发 WebSocket 服务不是只要一个框架就能跑得稳。关键在于你想清楚:当处理速度跟不上输入速度时,怎么办?背压就是你的安全阀。Actix Web 提供的基础设施很好,但如果你自己不设定边界,再牛的框架也没办法帮你挡住无底洞式的资源消耗。从有界队列、消息大小限制,到连接数上限、心跳清理,每一步都是为了一个目标:让系统在超负荷时优雅地减速,而不是直接崩溃。Rust 给了我们强力的编译器约束,但最终设计还是靠人。希望今天这些代码和思路,能让你以后再写 WebSocket 服务时,心里更有底。
评论
围绕“高并发WebSocket服务:基于Actix Web的背压控制与内存管理要点”参与讨论