第13章:综合实战 —— 从零实现多线程 Web 服务器
约 12 分钟 · 更新于 2026-09-01
第13章:综合实战 —— 从零实现多线程 Web 服务器
本章不把“浏览器能看到页面”当作终点,而是把一个同步 Web 服务拆成四个可验证的责任:接收连接、识别请求、限制并发、回收线程。我们会实现一个固定容量的工作队列,并用明确的关闭顺序保证已接收请求不会被悄悄丢弃。
所属教程
Rust 深度教程
导读:先写资源边界,再写业务分支
网络示例很容易从下面这几行开始:
rust
for connection in listener.incoming() {
handle(connection?);
}
它确实能工作,却没有回答三个决定服务稳定性的问题:
- 一个请求变慢时,后续连接在哪里等待?
- 突发流量到来时,服务器最多创建多少执行单元?
- 进程准备退出时,已经接收的请求由谁完成?
本章构建的教学项目叫 lantern-web,目标不是实现完整 HTTP,而是借一个足够小的协议表面观察并发架构。
text
lantern-web/
├── Cargo.toml
└── src/
├── main.rs # 监听、路由、关闭
└── executor.rs # 有界线程池
完成后,请求的流向如下:
text
TcpListener
│ accept
▼
主线程 ── Job ──▶ 有界队列 ──▶ 固定数量 Worker
│
├─ 读取请求头
├─ 选择路由
└─ 写回并关闭连接
协议边界
示例只处理单个 HTTP/1.1 请求,支持 GET /、GET /health 和 GET /pause,响应后主动关闭连接。它不实现 keep-alive、请求体、TLS、分块传输、压缩或完整的恶意输入防护。生产服务应使用 Hyper、Axum、Actix Web 等成熟组件。
一、定义验收标准
先把“做完”写成可以测试的结果:
- / 返回一段 HTML;
- /health 返回纯文本 ready;
- /pause 人为等待 1500 毫秒,用于制造慢请求;
- 慢请求不会阻止空闲 Worker 处理快请求;
- Worker 数固定,任务队列有容量上限;
- 队列已满时,提交方产生背压,而不是继续占用内存;
- 停止接收连接后,队列中的任务会被处理完;
- 主线程在所有 Worker 退出后再结束。
这些标准刻意把“功能正确”和“生命周期正确”放在一起。服务器是否可靠,不能只看状态码。
二、TCP 是字节流,HTTP 边界由应用识别
2.1 建立监听器
rust
use std::io;
use std::net::TcpListener;
fn main() -> io::Result<()> {
let listener = TcpListener::bind("127.0.0.1:7878")?;
println!("lantern-web listening at http://127.0.0.1:7878");
for accepted in listener.incoming() {
match accepted {
Ok(socket) => println!("new peer: {}", socket.peer_addr()?),
Err(error) => eprintln!("accept error: {error}"),
}
}
Ok(())
}
incoming() 的元素是 io::Result<TcpStream>,而不是永远成功的 TcpStream。连接建立会受到文件描述符、系统队列、网络状态和进程资源限制影响,因此主循环要把单次失败当作运行时事件,而不是“不可能发生”的分支。
2.2 一次 read 不等于一次请求
TCP 只承诺字节有序到达,不承诺应用消息的切片方式:
text
发送端:GET / HTTP/1.1\r\nHost: localhost\r\n\r\n
读取 A:GET / HT
读取 B:TP/1.1\r\nHost: loca
读取 C:lhost\r\n\r\n
HTTP 请求头用空行结束,也就是连续的 \r\n\r\n。教学实现使用 BufRead::read_line 按行累积,并设置总长度上限,避免客户端无限发送请求头。
三、把请求读取限制在一个小函数中
3.1 请求模型
本章只需要方法、路径和协议版本:
rust
#[derive(Debug, PartialEq, Eq)]
struct RequestLine {
method: String,
target: String,
version: String,
}
解析函数不做 I/O,便于单元测试:
rust
use std::io;
fn parse_request_line(line: &str) -> io::Result<RequestLine> {
let mut fields = line
.trim_end_matches(|ch| ch == '\r' || ch == '\n')
.split_whitespace();
let method = fields.next().ok_or_else(|| invalid("missing method"))?;
let target = fields.next().ok_or_else(|| invalid("missing target"))?;
let version = fields.next().ok_or_else(|| invalid("missing version"))?;
if fields.next().is_some() {
return Err(invalid("too many fields in request line"));
}
Ok(RequestLine {
method: method.to_owned(),
target: target.to_owned(),
version: version.to_owned(),
})
}
fn invalid(message: &'static str) -> io::Error {
io::Error::new(io::ErrorKind::InvalidData, message)
}
这里没有直接用字符串整体匹配 "GET / HTTP/1.1"。先形成结构,再做路由,可以独立区分“不支持的方法”“未知路径”和“协议格式错误”。
3.2 读取请求头
rust
use std::io::{self, BufRead, BufReader};
use std::net::TcpStream;
use std::time::Duration;
const MAX_HEADER_BYTES: usize = 8 * 1024;
fn read_request(stream: &mut TcpStream) -> io::Result<RequestLine> {
stream.set_read_timeout(Some(Duration::from_secs(5)))?;
let mut reader = BufReader::new(stream);
let mut total = 0usize;
let mut first = String::new();
let read = reader.read_line(&mut first)?;
if read == 0 {
return Err(invalid("peer closed before request line"));
}
total += read;
loop {
let mut header = String::new();
let read = reader.read_line(&mut header)?;
if read == 0 {
return Err(invalid("peer closed inside headers"));
}
total += read;
if total > MAX_HEADER_BYTES {
return Err(invalid("request headers are too large"));
}
if header == "\r\n" {
break;
}
}
parse_request_line(&first)
}
这段代码只检查总长度,没有逐项解释请求头。它仍然比“读取一个固定大小数组后直接转字符串”更接近真实协议处理,因为它明确区分了:
- 正常读完请求头;
- 对端提前关闭;
- 输入超过预算;
- socket 超时;
- 请求行格式非法。
为什么给 BufReader 传 &mut TcpStream
读取阶段只临时借用连接。函数返回后借用结束,调用方重新获得 TcpStream 的可变访问权,随后可以写响应。
四、建立响应模型,而不是到处拼字符串
rust
struct Response {
status: &'static str,
content_type: &'static str,
body: String,
}
impl Response {
fn html(status: &'static str, body: impl Into<String>) -> Self {
Self {
status,
content_type: "text/html; charset=utf-8",
body: body.into(),
}
}
fn text(status: &'static str, body: impl Into<String>) -> Self {
Self {
status,
content_type: "text/plain; charset=utf-8",
body: body.into(),
}
}
}
路由函数保持同步和纯粹,只有 /pause 负责模拟耗时:
rust
use std::thread;
use std::time::Duration;
fn route(request: &RequestLine) -> Response {
if request.method != "GET" {
return Response::text("405 Method Not Allowed", "only GET is supported");
}
match request.target.as_str() {
"/" => Response::html(
"200 OK",
"<!doctype html><meta charset=\"utf-8\"><h1>Lantern Web</h1>",
),
"/health" => Response::text("200 OK", "ready"),
"/pause" => {
thread::sleep(Duration::from_millis(1500));
Response::text("200 OK", "pause finished")
}
_ => Response::text("404 Not Found", "route not found"),
}
}
写响应时,Content-Length 使用字节数:
rust
use std::io::{self, Write};
use std::net::TcpStream;
fn write_response(stream: &mut TcpStream, response: Response) -> io::Result<()> {
let head = format!(
"HTTP/1.1 {}\r\nContent-Type: {}\r\nContent-Length: {}\r\nConnection: close\r\n\r\n",
response.status,
response.content_type,
response.body.as_bytes().len(),
);
stream.write_all(head.as_bytes())?;
stream.write_all(response.body.as_bytes())?;
stream.flush()
}
最后把一次连接的读取、计算和写入串起来:
rust
fn serve_one(mut stream: TcpStream) -> io::Result<()> {
let request = read_request(&mut stream)?;
let response = route(&request);
write_response(&mut stream, response)
}
五、先观察串行模型的真实症状
rust
for accepted in listener.incoming() {
match accepted {
Ok(stream) => {
if let Err(error) = serve_one(stream) {
eprintln!("request failed: {error}");
}
}
Err(error) => eprintln!("accept failed: {error}"),
}
}
打开两个终端,先请求慢路由,再立即请求健康检查:
bash
curl http://127.0.0.1:7878/pause
curl http://127.0.0.1:7878/health
第二个请求会等待第一个请求完成。此时并不是 /health 很慢,而是唯一执行线程被 /pause 占用,无法回到 accept。
text
时间 ───────────────────────────────────▶
主线程 accept → /pause 等待 1.5s → accept → /health
这个现象叫队头阻塞:排在前面的慢工作延迟了本可快速完成的后续工作。
六、每连接创建线程为什么只是过渡方案
最直接的并发写法如下:
rust
for accepted in listener.incoming() {
if let Ok(stream) = accepted {
std::thread::spawn(move || {
if let Err(error) = serve_one(stream) {
eprintln!("request failed: {error}");
}
});
}
}
move 把 TcpStream 交给新线程,因此新线程不借用主循环的局部变量。慢请求与快请求现在可以并行,但系统容量仍由外部连接数量决定:
text
连接数增长
├─ 线程栈占用增长
├─ 创建/销毁成本增长
├─ 调度切换增长
└─ 进程可用线程与文件描述符被耗尽
可靠的并发设计至少要显式给出两个数字:
因此下一步不是“更聪明地 spawn”,而是建立固定 Worker 数和有界队列。
七、设计一个有界执行器
7.1 公开接口
主程序只需要三个能力:
rust
let executor = BoundedExecutor::new(4, 16)?;
executor.submit(|| do_work())?;
drop(executor); // 排空队列并等待 Worker
- worker_count = 4:最多四个任务同时执行;
- queue_capacity = 16:最多十六个任务等待;
- 队列满时 submit 阻塞,背压传回 accept 线程;
- 执行器被销毁时不再接收新任务,Worker 处理完现有任务后退出。
7.2 任务类型
rust
type Job = Box<dyn FnOnce() + Send + 'static>;
| 约束 | 表达的事实 |
|---|
| FnOnce() | 任务只会调用一次,并允许消费捕获值 |
| Send | 任务可以从提交线程移动到 Worker |
| 'static | 任务不依赖可能提前失效的栈借用 |
| Box<dyn ...> | 不同闭包被统一为固定大小的队列元素 |
'static 并不表示任务永不结束。它表示闭包拥有的数据足以独立存活,例如 String、TcpStream、Arc<T>。
7.3 数据结构
rust
use std::fmt;
use std::sync::{mpsc, Arc, Mutex};
use std::thread;
pub struct BoundedExecutor {
sender: Option<mpsc::SyncSender<Job>>,
workers: Vec<Worker>,
}
struct Worker {
name: String,
handle: Option<thread::JoinHandle<()>>,
}
#[derive(Debug)]
pub enum BuildError {
ZeroWorkers,
ZeroQueue,
}
impl fmt::Display for BuildError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
match self {
Self::ZeroWorkers => write!(f, "worker count must be greater than zero"),
Self::ZeroQueue => write!(f, "queue capacity must be greater than zero"),
}
}
}
impl std::error::Error for BuildError {}
SyncSender 来自 sync_channel。当缓冲区已满,send 会等待队列腾出位置,这使过载可见,而不是把任务无限堆在内存里。
八、实现 Worker:锁只覆盖领取动作
标准库 channel 只有一个 Receiver,多个 Worker 通过 Arc<Mutex<_>> 共享它:
rust
impl Worker {
fn spawn(
index: usize,
receiver: Arc<Mutex<mpsc::Receiver<Job>>>,
) -> std::io::Result<Self> {
let name = format!("lantern-worker-{index}");
let thread_name = name.clone();
let handle = thread::Builder::new()
.name(thread_name.clone())
.spawn(move || loop {
let received = {
let guard = match receiver.lock() {
Ok(guard) => guard,
Err(poisoned) => {
eprintln!("{thread_name}: receiver lock was poisoned");
poisoned.into_inner()
}
};
guard.recv()
};
let job = match received {
Ok(job) => job,
Err(_) => break,
};
let outcome = std::panic::catch_unwind(
std::panic::AssertUnwindSafe(job),
);
if outcome.is_err() {
eprintln!("{thread_name}: job panicked; worker remains available");
}
})?;
Ok(Self {
name,
handle: Some(handle),
})
}
}
关键是这段显式作用域:
rust
let received = {
let guard = receiver.lock().unwrap();
guard.recv()
}; // MutexGuard 在这里释放
let job = received?;
job();
如果执行 job() 时仍持有接收锁,其他 Worker 无法领取任务,线程池会在表面并发、实际串行的状态下运行。
不要用代码长度推断锁生命周期
并发代码应让锁的结束位置一眼可见。尤其要谨慎对待把 lock()、recv() 和循环条件压在同一表达式里的写法。
catch_unwind 让单个任务 panic 后 Worker 仍可继续服务。它不是所有错误的解决方案:如果任务修改了共享状态,仍需判断该状态在 panic 后是否有效。
九、构造、提交与关闭
9.1 构造执行器
rust
impl BoundedExecutor {
pub fn new(worker_count: usize, queue_capacity: usize) -> Result<Self, BuildError> {
if worker_count == 0 {
return Err(BuildError::ZeroWorkers);
}
if queue_capacity == 0 {
return Err(BuildError::ZeroQueue);
}
let (sender, receiver) = mpsc::sync_channel(queue_capacity);
let receiver = Arc::new(Mutex::new(receiver));
let mut workers = Vec::with_capacity(worker_count);
for index in 0..worker_count {
let worker = Worker::spawn(index, Arc::clone(&receiver))
.expect("failed to create worker thread");
workers.push(worker);
}
Ok(Self {
sender: Some(sender),
workers,
})
}
}
教学代码把线程创建失败保留为 expect,是为了让 BuildError 聚焦参数错误。进一步工程化时,应把 std::io::Error 纳入构造错误,并处理“部分 Worker 已创建、后续创建失败”的清理问题。
9.2 提交任务
rust
#[derive(Debug)]
pub struct SubmitError;
impl fmt::Display for SubmitError {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
write!(f, "executor is no longer accepting jobs")
}
}
impl std::error::Error for SubmitError {}
impl BoundedExecutor {
pub fn submit<F>(&self, task: F) -> Result<(), SubmitError>
where
F: FnOnce() + Send + 'static,
{
let sender = self.sender.as_ref().ok_or(SubmitError)?;
sender.send(Box::new(task)).map_err(|_| SubmitError)
}
}
当队列满时,send 会阻塞 accept 线程。这是一种简单背压:操作系统连接 backlog 与应用任务队列共同承担等待。它不一定适合所有服务,但比无界增长更容易推理。
其他可选策略包括:
- try_send 立即拒绝,并返回 503;
- 等待固定时长后拒绝;
- 按路由设置不同容量;
- 在进入线程池前使用连接限流器。
策略没有统一答案,重要的是让过载行为成为 API 的明确部分。
9.3 用 channel 断开表达关闭
销毁执行器时,先释放最后一个发送端。Worker 处理完队列中的任务后,recv() 返回错误,循环自然退出:
rust
impl Drop for BoundedExecutor {
fn drop(&mut self) {
drop(self.sender.take());
for worker in &mut self.workers {
if let Some(handle) = worker.handle.take() {
if handle.join().is_err() {
eprintln!("{} terminated with panic", worker.name);
}
}
}
}
}
关闭顺序不能反过来。如果先 join,Worker 仍可能阻塞在 recv(),主线程会永久等待。
JoinHandle::join(self) 消费句柄,而 Drop::drop 只有 &mut self。把句柄放在 Option 中,便可通过 take() 将所有权移出字段。
text
Running
│ 停止 accept
▼
NoNewConnections
│ drop Sender
▼
DrainingQueue
│ recv 返回断开
▼
WorkersExit
│ join 全部句柄
▼
Closed
十、组合完整服务器
src/main.rs 的核心结构如下:
rust
mod executor;
use executor::BoundedExecutor;
use std::error::Error;
use std::net::TcpListener;
fn main() -> Result<(), Box<dyn Error>> {
let listener = TcpListener::bind("127.0.0.1:7878")?;
let executor = BoundedExecutor::new(4, 16)?;
println!("listening at http://127.0.0.1:7878");
// take(6) 只用于稳定演示关闭流程。
for accepted in listener.incoming().take(6) {
match accepted {
Ok(stream) => {
executor.submit(move || {
if let Err(error) = serve_one(stream) {
eprintln!("connection failed: {error}");
}
})?;
}
Err(error) => eprintln!("accept failed: {error}"),
}
}
println!("accept loop stopped; draining queued requests");
drop(executor);
println!("all workers joined");
Ok(())
}
启动服务后,可并行发起请求:
另开终端:
bash
curl http://127.0.0.1:7878/pause &
curl http://127.0.0.1:7878/health &
curl http://127.0.0.1:7878/ &
wait
如果有空闲 Worker,/health 和 / 不必等待 /pause。当累计接收六个连接后,主循环停止;执行器先处理已入队任务,再回收线程。
真实服务不会使用 take(6),而会监听 Ctrl+C、SIGTERM、管理命令或编排平台的终止信号。无论触发源是什么,关闭顺序仍应保持:停止入口、排空或拒绝、通知退出、等待回收。
十一、如何验证,而不只靠肉眼观察
11.1 请求行解析测试
rust
#[test]
fn parses_request_line() {
let line = parse_request_line("GET /health HTTP/1.1\r\n").unwrap();
assert_eq!(line.method, "GET");
assert_eq!(line.target, "/health");
assert_eq!(line.version, "HTTP/1.1");
}
#[test]
fn rejects_extra_fields() {
assert!(parse_request_line("GET / HTTP/1.1 unexpected\r\n").is_err());
}
11.2 验证任务确实并发
可使用 barrier 让多个任务同时进入临界观察点:
rust
use std::sync::{Arc, Barrier};
#[test]
fn runs_two_jobs_at_the_same_time() {
let executor = BoundedExecutor::new(2, 2).unwrap();
let gate = Arc::new(Barrier::new(3));
for _ in 0..2 {
let gate = Arc::clone(&gate);
executor.submit(move || {
gate.wait();
}).unwrap();
}
gate.wait();
}
如果只有一个 Worker,这个测试会卡住,因此测试环境还应配合超时机制,避免失败时永久挂起。
11.3 验证 Drop 等待在途任务
rust
use std::sync::mpsc;
#[test]
fn drop_waits_for_queued_work() {
let (done_tx, done_rx) = mpsc::channel();
{
let executor = BoundedExecutor::new(1, 1).unwrap();
executor.submit(move || done_tx.send("finished").unwrap()).unwrap();
}
assert_eq!(done_rx.recv().unwrap(), "finished");
}
测试的重点不是线程池内部字段,而是公共生命周期合同:作用域结束后,已接受任务已经完成。
十二、设计取舍复盘
| 设计点 | 本章选择 | 得到的性质 | 仍然存在的限制 |
|---|
| I/O 模型 | 阻塞 socket | 代码路径直观 | 每个活跃请求占据 Worker |
| 执行模型 | 固定线程池 | 并发数可控 | 不适合海量空闲长连接 |
| 等待队列 | sync_channel | 容量有上限 | 队列满会阻塞 accept 线程 |
| 共享接收端 | Arc<Mutex<Receiver>> | 仅用标准库实现 | 领取动作被串行化 |
| panic 隔离 | catch_unwind | Worker 不因单个任务消失 | 共享业务状态仍需恢复策略 |
| 关闭信号 | drop 最后一个 Sender | 无额外布尔状态 | 必须确保没有遗留 Sender clone |
| HTTP 范围 | 单请求、主动关闭 | 教学边界清晰 | 不具备生产协议完整性 |
线程池首先是容量治理工具,其次才是并发加速工具。固定 Worker 限制“正在执行”,有界队列限制“尚未执行”,关闭协议限制“生命周期结束时允许遗留什么”。
十三、常见失误与诊断线索
13.1 把网络错误全部 unwrap
症状:一个客户端断开导致整个进程退出。
处理:在 bind、accept、单连接处理、任务提交、关闭阶段分别记录错误,让局部失败停留在局部。
13.2 Worker 执行任务时仍持有接收锁
症状:Worker 数大于一,但 /pause 仍阻塞所有请求。
处理:用独立作用域取得 Job,确认 MutexGuard 在调用任务前销毁。
13.3 队列没有容量
症状:压测时吞吐不再增长,内存和尾延迟却持续增加。
处理:使用有界队列,定义满载时等待、超时或拒绝策略,并记录队列深度。
13.4 关闭时先等待线程
症状:进程收到关闭指令后不退出,线程栈停在 recv()。
处理:先关闭所有任务发送入口,再 join。
13.5 只限制 Worker,不限制连接读取
症状:慢客户端长时间占据所有 Worker,健康检查也无法执行。
处理:设置读取超时、请求头上限,必要时把连接接入与请求执行拆成不同层级。
13.6 把“固定线程池”理解为完整过载保护
症状:操作系统 backlog、应用队列和下游依赖仍被压满。
处理:从入口到下游逐段标注容量,避免只在某一层设置一个数字。
十四、练习
练习 1:拒绝而不是阻塞
为执行器增加 try_submit。队列满时把连接返回给调用方,并写回 503 Service Unavailable。比较它与阻塞提交在延迟和吞吐上的差异。
练习 2:分级队列
让 /health 不与 /pause 共用同一队列。设计两组 Worker 或一个带优先级的调度器,并说明如何避免低优先级任务永久饥饿。
练习 3:部分构造失败
让 BoundedExecutor::new 返回线程创建错误。模拟第三个 Worker 创建失败,验证前两个 Worker 能被正确关闭和回收。
练习 4:关闭截止时间
为关闭流程增加最大等待时长。超过期限后输出仍未完成的任务数量,并讨论标准线程无法安全强杀时有哪些替代方案。
练习 5:输入防护
增加方法长度、路径长度、请求头数量和单行长度限制;为超限情况分别返回 400、414 或 431。
练习 6:可观测性
记录请求 ID、路由、排队时间、执行时间、Worker 名称和结果。确保日志不包含敏感请求头。
练习 7:对照异步实现
用 Tokio 重写监听与连接任务,保持路由和响应模型不变。比较:
- 阻塞线程与异步 task 的资源占用;
- 队列背压的位置;
- 超时与取消实现;
- 优雅关闭的等待方式。
十五、总结
- TCP 提供连续字节,HTTP 消息边界必须由协议解析器识别。
- 解析、路由和写响应分层后,每一层都可以独立测试。
- 单线程服务器的主要问题是慢请求造成队头阻塞。
- 每连接创建线程能获得并发,却把外部负载直接转换成无界线程资源。
- 固定 Worker 数控制执行并发,有界任务队列控制等待规模。
- FnOnce + Send + 'static 描述一次性、可跨线程且不依赖短期借用的任务。
- Arc<Mutex<Receiver<_>>> 只应用于领取任务;锁不能覆盖任务执行。
- 队列满时如何处理是服务合同,不是可以忽略的实现细节。
- 释放最后一个 Sender 可以同时表达“停止提交”和“队列排空后退出”。
- Option<JoinHandle> 让 Drop 能通过 take() 取得句柄所有权。
- 优雅关闭的顺序是停止入口、处理遗留、通知退出、等待回收。
- 教学 HTTP 服务帮助理解资源边界,但不能替代成熟协议栈。
章节导航
- 上一章:第12章:Cargo 工程化
- 下一章:第14章:用 Tokio 实现 Mini-Redis
- Rust 深度教程
- 第3章:所有权与借用
- 第7章:闭包与迭代器
- 第9章:并发编程
- 第12章:Cargo 工程化
- 第14章:Tokio Mini-Redis