阶段四 · 并发编排与 IO

实战:异步 TCP 服务器

一句话总结

把前面所有零件组装起来:TcpListener 循环接受连接,每来一个连接就 spawn 一个任务去处理--这就是"接受循环 + 每连接一任务"模式,异步服务器的标准骨架.
本篇从一个 echo 服务器起步,再升级成多客户端广播聊天室.

本篇目标

这是阶段四的综合实战.我们要写一个能同时服务大量客户端的 TCP 服务器,用到的全是前面学过的东西:

  • spawn(第 06 篇)--每个连接一个独立任务
  • 异步读写 AsyncReadExt/AsyncWriteExt(第 11 篇)
  • broadcast channel(第 09 篇)--实现聊天广播
  • select!(第 10 篇)--同时处理"收网络数据"与"收广播消息"

第一步:echo 服务器

echo(回声)服务器把客户端发来的数据原样发回.它最小地展示了服务器骨架:

use tokio::net::TcpListener;
use tokio::io::{AsyncReadExt, AsyncWriteExt};

#[tokio::main]
async fn main() -> std::io::Result<()> {
let listener = TcpListener::bind("127.0.0.1:8080").await?;
println!("监听 127.0.0.1:8080");

loop {
let (mut socket, addr) = listener.accept().await?;
println!("新连接:{addr}");

tokio::spawn(async move {
let mut buf = [0u8; 1024];
loop {
match socket.read(&mut buf).await {
Ok(0) => return, // 读到 0 字节 = 对端关闭
Ok(n) => {
if socket.write_all(&buf[..n]).await.is_err() {
return;
}
}
Err(_) => return,
}
}
});
}
}

为什么必须 spawn,不能直接处理

如果把 handle 逻辑直接写在接受循环里(不 spawn),那么处理当前连接的整个过程会阻塞接受循环,服务器同一时刻只能服务一个客户端--退化成串行.

spawn 让每个连接的处理独立推进,接受循环立即回头等下一个连接.这就是"每连接一任务"的精髓.

接受循环像餐厅领位员:只负责把客人迎进门,安排一位服务员(spawn 一个任务)跟进,然后马上回到门口迎下一位.领位员从不亲自上菜--否则门口会排起长队.

动手验证

跑起来后,用系统自带的 nc(netcat)就能测:

# 终端 A:启动服务器
cargo run

# 终端 B:连接并输入,会看到输入被原样返回
nc 127.0.0.1 8080
hello # 输入
hello # 服务器回显

# 再开终端 C 同时连接,验证多客户端并发(互不影响)
nc 127.0.0.1 8080

第二步:升级为广播聊天室

echo 只是自己跟自己玩.真正体现并发价值的是多客户端互通:任一客户端发的消息,广播给所有在线客户端.这需要每个连接任务同时做两件事:

  1. 读本连接客户端发来的数据 → 广播出去
  2. 收别人广播来的消息 → 写给本连接客户端

"同时等两件事,哪个先来处理哪个"--正是 select! 的舞台,而广播用 broadcast channel:

use tokio::net::TcpListener;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::sync::broadcast;

#[tokio::main]
async fn main() -> std::io::Result<()> {
let listener = TcpListener::bind("127.0.0.1:8080").await?;
let (tx, _rx) = broadcast::channel::<String>(16);

loop {
let (mut socket, addr) = listener.accept().await?;
let tx = tx.clone();
let mut rx = tx.subscribe();

tokio::spawn(async move {
let mut buf = [0u8; 1024];
loop {
tokio::select! {
result = socket.read(&mut buf) => {
match result {
Ok(0) | Err(_) => break,
Ok(n) => {
let msg = format!("{addr}: {}",
String::from_utf8_lossy(&buf[..n]));
let _ = tx.send(msg);
}
}
}
Ok(msg) = rx.recv() => {
if socket.write_all(msg.as_bytes()).await.is_err() {
break;
}
}
}
}
});
}
}

把零件对号入座

  • broadcast:一条消息要送达所有客户端 → 正是第 09 篇 broadcast 的典型场景.
  • select!:每个连接任务要同时盯"网络输入"和"广播消息"两个 Future,谁先就绪处理谁 → 第 10 篇.
  • spawn:每个客户端一个任务,几千个连接也只是几千个轻量任务 → 第 06 篇.

这里的 select! 取消安全吗

回顾第 10 篇的警告:select! 未选中的分支会被取消.这里 socket.readrx.recv 都是取消安全的--被取消时不会丢失已确认的数据(read 未读走的字节仍在内核缓冲,recv 未取走的消息仍在信道).

所以这个循环是安全的.换成非取消安全的操作(如 read_exact 读到一半)就可能出问题.

为什么它能扛住海量连接

回到第 01 篇的初心.这个服务器若用"一连接一线程"写,1 万个客户端就是 1 万个线程,约 20 GB 栈.而这里 1 万个连接只是 1 万个异步任务,跑在 Tokio 的少数几个 worker 线程上,内存占用低一两个数量级.绝大部分连接在"等输入"时不占用任何 CPU--这正是异步的全部意义在一个真实程序里的兑现.

快速回顾

  • 标准骨架:TcpListener 接受循环 + 每连接 spawn 一个任务;接受循环绝不亲自处理连接.
  • echo 服务器:read → write_all 回写;读到 0 字节即对端关闭.
  • 广播聊天室:broadcast 分发消息 + select! 同时处理"读输入"与"收广播".
  • 取消安全:select! 循环里的 read/recv 取消安全,故可放心被反复取消重建.
  • 规模兑现:万级连接 = 万级轻量任务而非万级线程,内存与调度开销低一两个数量级.

动手练习

  1. 验证并发:跑通 echo 服务器,用两个 nc 同时连接,确认互不阻塞.
  2. 退化实验:故意把处理逻辑写在接受循环里(不 spawn),用两个客户端验证它退化为串行.
  3. 广播聊天:跑通广播聊天室,开三个 nc 客户端,确认一个发言其余都能收到.
  4. 扩展功能:给聊天室加一条:新客户端连入时,向所有人广播"某地址加入了".
  5. Lagged 处理:思考并尝试:如果某客户端读得很慢导致 broadcast Lagged,该如何处理?(结合第 09 篇)