Agent X-Ray
RuntimeNotesAbout
Notes/代码工程/Rust 深度教材/第14章

第14章:异步实战 —— 用 Tokio 实现 Mini-Redis

17 分钟 · 更新于 2026-09-01

第14章:异步实战 —— 用 Tokio 实现 Mini-Redis

本章构建一个名为 PocketKV 的教学型键值服务。我们不追求复刻 Redis 的全部行为,而是沿着“连接任务—协议帧—命令—状态—关闭监督”这条主线,练习 Tokio 中最重要的工程判断:什么可以并发,什么必须只有一个所有者,在哪里设置容量,以及 Future 被取消后系统处于什么状态。

所属教程 Rust 深度教程


导读:异步不是把线程换成 .await

第13章的固定线程池适合数量受控的阻塞工作,但一个连接在等待读取时会独占 Worker。对于大量“多数时间在等网络”的连接,更合适的模型是:把每条连接表示成 Future,在 socket 未就绪时把执行权交还运行时。

text
操作系统就绪事件
      │
      ▼
Tokio Runtime ──唤醒──▶ 可运行 Task
      │                    │
      ├─ socket read       ├─ 解析 Frame
      ├─ timer             ├─ 执行 Command
      └─ signal            └─ 写回 Response

这并不意味着资源自动有上限。每条 task 仍会持有连接、缓冲区、状态句柄和排队消息。本章从一开始就把容量、取消和关闭纳入设计。

最终项目结构:

text
pocket-kv/
├── Cargo.toml
└── src/
    ├── bin/
    │   ├── server.rs
    │   └── client.rs
    ├── client.rs
    ├── connection.rs
    ├── error.rs
    ├── frame.rs
    ├── store.rs
    └── lib.rs

实现范围:

  • PING [message]
  • GET key
  • SET key value
  • 基于简化 RESP 的 Frame 编解码
  • 每连接一个 task,并限制最大连接数
  • 共享存储与 manager task 两种状态组织方式
  • 有界 mpsc + oneshot 客户端句柄
  • select! 超时与关闭
  • Stream 形式的事件订阅
  • JoinSet 监督和带截止时间的优雅关闭
  • 同步调用方复用异步客户端

心智模型 调用 async fn 只会创建 Future;.await 让当前 Future 运行到完成或暂时无法继续;tokio::spawn 才把 Future 注册为独立 task。异步提升的是等待期间的复用能力,不会自动修复无界队列、锁竞争或慢下游。


一、项目依赖与 feature 预算

bash
cargo new pocket-kv
cd pocket-kv

Cargo.toml

toml
[package]
name = "pocket-kv"
version = "0.1.0"
edition = "2021"

[dependencies]
bytes = "1"
thiserror = "2"
tokio = {
    version = "1",
    features = [
        "io-util",
        "macros",
        "net",
        "rt-multi-thread",
        "signal",
        "sync",
        "time",
    ],
}
tokio-stream = { version = "0.1", features = ["sync"] }

直接开启 Tokio 的 full feature 很方便,但教程工程更值得练习按能力选择 feature。这样做可以减少编译面,也迫使我们理解代码真正依赖了 runtime 的哪些组成部分。

参考:第12章:Cargo 工程化。


二、把连接接入变成受控并发

2.1 串行写法仍然是串行

rust
use tokio::net::{TcpListener, TcpStream};

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

    loop {
        let (socket, peer) = listener.accept().await?;
        println!("accepted {peer}");
        serve_connection(socket).await?;
    }
}

async fn serve_connection(_socket: TcpStream) -> std::io::Result<()> {
    Ok(())
}

虽然函数带有 async,但 serve_connection(...).await 完成前,循环不会执行下一次 accept。异步语法本身不等于并发。

2.2 spawn 建立独立连接任务

rust
loop {
    let (socket, peer) = listener.accept().await?;

    tokio::spawn(async move {
        if let Err(error) = serve_connection(socket).await {
            eprintln!("peer {peer} failed: {error}");
        }
    });
}

async movesocketpeer 的所有权移入 Future。多线程 runtime 上的 spawned Future 通常需要 Send + 'static

  • Send:task 挂起后可能在另一个 Worker 线程继续;
  • 'static:task 不能借用即将离开作用域的局部数据。

'static 约束生命周期依赖,不要求 task 永久运行。

2.3 用 Semaphore 限制连接数

轻量 task 不是零成本 task。服务器接受连接前先取得 permit:

rust
use std::sync::Arc;
use tokio::sync::Semaphore;
use tokio::task::JoinSet;

let permits = Arc::new(Semaphore::new(512));
let mut connections = JoinSet::new();

loop {
    let permit = Arc::clone(&permits)
        .acquire_owned()
        .await
        .expect("semaphore closed unexpectedly");

    let (socket, peer) = listener.accept().await?;

    connections.spawn(async move {
        let _permit = permit;
        let result = serve_connection(socket).await;
        if let Err(error) = &result {
            eprintln!("peer {peer} failed: {error}");
        }
        result
    });
}

permit 的所有权跟随 task;task 结束时自动释放。这个写法会在没有 permit 时暂停继续 accept。如果希望立即接收再返回“服务繁忙”,可以改用 try_acquire_owned,但要先考虑操作系统 backlog 与拒绝响应的成本。

容量应该贴近被保护的资源 连接上限保护 socket、task 与缓冲区;命令队列容量保护 manager;存储上限保护内存。一个全局数字无法替代逐层容量设计。


三、协议分层:字节、Frame、Command 各司其职

3.1 Frame 只描述线格式

rust
use bytes::Bytes;

#[derive(Debug, Clone, PartialEq, Eq)]
pub enum Frame {
    Simple(String),
    Error(String),
    Integer(i64),
    Bulk(Bytes),
    Null,
    Array(Vec<Frame>),
}

Frame 不知道 GET 的含义,也不访问数据库。层次关系保持单向:

text
TcpStream
   ▼
Connection:缓冲读写
   ▼
Frame:协议语法
   ▼
Command:参数与业务语义
   ▼
Store:状态变更

如果 parser 直接操作 HashMap,畸形输入测试就会被网络和数据库状态绑在一起,协议升级也更难。

3.2 Connection 保存未消费字节

rust
use bytes::BytesMut;
use tokio::io::BufWriter;
use tokio::net::TcpStream;

pub struct Connection {
    stream: BufWriter<TcpStream>,
    input: BytesMut,
}

impl Connection {
    pub fn new(stream: TcpStream) -> Self {
        Self {
            stream: BufWriter::new(stream),
            input: BytesMut::with_capacity(4 * 1024),
        }
    }
}

一次 read 可能得到:

  • 半个 Frame;
  • 恰好一个 Frame;
  • 一个完整 Frame 加下一个 Frame 的前缀;
  • 多个连续 Frame。

因此 input 必须跨多次读取保留数据。

3.3 先解析缓存,再读 socket

rust
use tokio::io::AsyncReadExt;

const MAX_BUFFERED_BYTES: usize = 1024 * 1024;

impl Connection {
    pub async fn read_frame(&mut self) -> Result<Option<Frame>, ProtocolError> {
        loop {
            if let Some(frame) = self.try_parse()? {
                return Ok(Some(frame));
            }

            if self.input.len() >= MAX_BUFFERED_BYTES {
                return Err(ProtocolError::FrameTooLarge);
            }

            let read = self.stream.read_buf(&mut self.input).await?;
            if read == 0 {
                return if self.input.is_empty() {
                    Ok(None)
                } else {
                    Err(ProtocolError::TruncatedFrame)
                };
            }
        }
    }
}

顺序不能写反。缓存中可能已经包含完整 Frame,先读网络会造成不必要等待。EOF 也要分成两类:空缓存表示连接在帧边界正常结束,非空缓存表示对端留下了半帧。

3.4 检查长度与构造对象分开

rust
use bytes::Buf;
use std::io::Cursor;

impl Connection {
    fn try_parse(&mut self) -> Result<Option<Frame>, ProtocolError> {
        let mut cursor = Cursor::new(&self.input[..]);

        match Frame::check(&mut cursor) {
            Ok(()) => {
                let consumed = cursor.position() as usize;
                cursor.set_position(0);

                let frame = Frame::parse(&mut cursor)?;
                self.input.advance(consumed);
                Ok(Some(frame))
            }
            Err(ProtocolError::Incomplete) => Ok(None),
            Err(error) => Err(error),
        }
    }
}

check 只遍历结构并判断完整长度,parse 在确认数据完整后再分配 StringBytes 和数组。这种两阶段策略让“不完整”成为正常状态,而不是异常恢复路径。

生产解析器还应限制:

  • Bulk 长度;
  • Array 元素数量;
  • 递归深度;
  • 单帧读取时间;
  • 单连接累计缓冲;
  • 数字字段的溢出。

3.5 写 Frame

rust
use tokio::io::AsyncWriteExt;

impl Connection {
    pub async fn write_frame(&mut self, frame: &Frame) -> std::io::Result<()> {
        match frame {
            Frame::Simple(text) => {
                self.stream.write_u8(b'+').await?;
                self.stream.write_all(text.as_bytes()).await?;
                self.stream.write_all(b"\r\n").await?;
            }
            Frame::Bulk(bytes) => {
                self.stream.write_u8(b'$').await?;
                self.stream
                    .write_all(bytes.len().to_string().as_bytes())
                    .await?;
                self.stream.write_all(b"\r\n").await?;
                self.stream.write_all(bytes).await?;
                self.stream.write_all(b"\r\n").await?;
            }
            Frame::Null => self.stream.write_all(b"$-1\r\n").await?,
            Frame::Error(message) => {
                self.stream.write_u8(b'-').await?;
                self.stream.write_all(message.as_bytes()).await?;
                self.stream.write_all(b"\r\n").await?;
            }
            _ => return Err(std::io::Error::other("frame writer not implemented")),
        }

        self.stream.flush().await
    }
}

BufWriter 合并细碎写入,减少系统调用。每帧立即 flush 是延迟优先的教学选择;是否批量 flush 应通过真实负载基准决定。


四、命令类型把合法输入收窄

rust
use bytes::Bytes;

#[derive(Debug, PartialEq, Eq)]
enum Command {
    Ping(Option<Bytes>),
    Get { key: String },
    Set { key: String, value: Bytes },
}

从 Frame 到 Command 的转换应是纯函数:

rust
impl TryFrom<Frame> for Command {
    type Error = CommandError;

    fn try_from(frame: Frame) -> Result<Self, Self::Error> {
        let Frame::Array(mut parts) = frame else {
            return Err(CommandError::ExpectedArray);
        };

        let name = take_text(&mut parts)?.to_ascii_uppercase();

        match name.as_str() {
            "PING" => parse_ping(parts),
            "GET" => parse_get(parts),
            "SET" => parse_set(parts),
            _ => Err(CommandError::Unknown(name)),
        }
    }
}

解析阶段负责:

  • 命令名大小写归一化;
  • 参数数量;
  • 参数类型;
  • UTF-8 要求;
  • 多余参数;
  • 空 key 等业务前置条件。

执行阶段不再处理“第三个元素是不是 Bulk”之类的协议细节,只匹配合法 Command。


五、共享状态方案一:短临界区的同步 Mutex

5.1 Store 封装锁

rust
use bytes::Bytes;
use std::collections::HashMap;
use std::sync::Mutex;

#[derive(Default)]
pub struct Store {
    values: Mutex<HashMap<String, Bytes>>,
}

impl Store {
    pub fn get(&self, key: &str) -> Option<Bytes> {
        self.values
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner())
            .get(key)
            .cloned()
    }

    pub fn set(&self, key: String, value: Bytes) {
        self.values
            .lock()
            .unwrap_or_else(|poisoned| poisoned.into_inner())
            .insert(key, value);
    }
}

Bytes::clone 通常只增加引用计数,适合在锁内取出句柄后尽快释放锁。

5.2 为什么 async 代码可以使用 std::sync::Mutex

判断标准不是“函数是不是 async”,而是临界区是否满足:

  • 很短;
  • 不执行 I/O;
  • 不等待 timer 或 channel;
  • 不跨 .await
  • 竞争量可接受。

错误写法:

rust
let mut guard = store.values.lock().unwrap();
guard.insert(key, value);
connection.write_frame(&Frame::Simple("OK".into())).await?;

task 在写 socket 时可能挂起,锁却继续被持有。正确做法是让状态操作在同步方法中完成,再进入 .await

rust
store.set(key, value);
connection.write_frame(&Frame::Simple("OK".into())).await?;

把锁字段设为私有,可以从模块边界上降低误用概率。

5.3 执行命令

rust
fn execute(command: Command, store: &Store) -> Frame {
    match command {
        Command::Ping(message) => message
            .map(Frame::Bulk)
            .unwrap_or_else(|| Frame::Simple("PONG".into())),
        Command::Get { key } => store
            .get(&key)
            .map(Frame::Bulk)
            .unwrap_or(Frame::Null),
        Command::Set { key, value } => {
            store.set(key, value);
            Frame::Simple("OK".into())
        }
    }
}

执行函数保持同步,说明当前命令没有理由跨 .await 持有数据库状态。未来加入磁盘、远程服务或事务时,应重新设计边界,而不是把 .await 硬塞进锁作用域。


六、共享状态方案二:单一所有者与消息协议

当状态需要严格顺序、复杂事务或异步操作时,可以让一个 manager task 独占资源。其他 task 不拿锁,只发命令。

text
调用方 Task ──有界 mpsc──▶ Manager Task ──唯一拥有 Store/Client
调用方 Task ◀──oneshot──── Manager Task

6.1 定义管理请求

rust
use tokio::sync::oneshot;

#[derive(Debug)]
enum StoreRequest {
    Get {
        key: String,
        reply: oneshot::Sender<Option<Bytes>>,
    },
    Set {
        key: String,
        value: Bytes,
        reply: oneshot::Sender<()>,
    },
}

6.2 有界命令队列

rust
use tokio::sync::mpsc;

let (request_tx, mut request_rx) = mpsc::channel::<StoreRequest>(128);

容量 128 是协议的一部分:

  • 未满:发送进入队列;
  • 已满:send(...).await 暂停调用 task;
  • Receiver 结束:后续发送失败;
  • 所有 Sender 释放:manager 的 recv().await 返回 None

无界队列会把处理能力不足隐藏成越来越长的延迟和越来越高的内存占用。

6.3 Manager 循环

rust
let manager = tokio::spawn(async move {
    let mut values = HashMap::<String, Bytes>::new();

    while let Some(request) = request_rx.recv().await {
        match request {
            StoreRequest::Get { key, reply } => {
                let value = values.get(&key).cloned();
                let _ = reply.send(value);
            }
            StoreRequest::Set { key, value, reply } => {
                values.insert(key, value);
                let _ = reply.send(());
            }
        }
    }
});

oneshot::Sender::send 不需要 .await。接收方可能已因超时或取消而离开,所以发送失败通常只表示“结果无人接收”。

6.4 可克隆句柄

rust
#[derive(Clone)]
pub struct StoreHandle {
    tx: mpsc::Sender<StoreRequest>,
}

impl StoreHandle {
    pub async fn get(&self, key: impl Into<String>) -> Result<Option<Bytes>, StoreError> {
        let (reply_tx, reply_rx) = oneshot::channel();

        self.tx
            .send(StoreRequest::Get {
                key: key.into(),
                reply: reply_tx,
            })
            .await
            .map_err(|_| StoreError::Closed)?;

        reply_rx.await.map_err(|_| StoreError::Closed)
    }
}

这里有两次生命周期检查:请求能否进入 manager 队列,以及 manager 能否在退出前给出响应。把两处都写成 unwrap 会让关闭与故障难以区分。

选择锁还是 manager 的判断方法 如果操作是短小同步变更,锁往往更直接;如果资源必须单线程拥有、操作需要严格排序、要统一做批处理或重连,manager task 更容易表达不变量。消息传递不是天然更快,它的主要收益是所有权清晰。


七、连接循环:将 I/O、解析、执行串起来

rust
use std::sync::Arc;
use tokio::net::TcpStream;
use tokio::sync::watch;

async fn serve_connection(
    socket: TcpStream,
    store: Arc<Store>,
    mut shutdown: watch::Receiver<bool>,
) -> Result<(), ServerError> {
    let mut connection = Connection::new(socket);

    loop {
        tokio::select! {
            frame = connection.read_frame() => {
                let Some(frame) = frame? else {
                    return Ok(());
                };

                let response = match Command::try_from(frame) {
                    Ok(command) => execute(command, &store),
                    Err(error) => Frame::Error(error.to_string()),
                };

                connection.write_frame(&response).await?;
            }
            changed = shutdown.changed() => {
                if changed.is_err() || *shutdown.borrow() {
                    return Ok(());
                }
            }
        }
    }
}

非法命令被转换为 Error Frame,而不是 panic 或杀死整个服务。网络错误则结束当前连接,不影响其他 task。

这段代码还隐含一个取消选择:关闭信号赢得 select! 时,正在等待的 read_frame Future 被 drop。对“等待下一帧”而言通常可以接受,因为随后连接整体被关闭;但如果正在执行“写到一半的响应”或数据库事务,就必须先定义是否允许取消。


八、select! 是取消边界,不只是语法糖

8.1 多路等待

rust
use tokio::time::{sleep, Duration};

tokio::select! {
    result = connection.read_frame() => {
        handle(result?).await?;
    }
    _ = sleep(Duration::from_secs(30)) => {
        return Err(ServerError::IdleTimeout);
    }
    _ = shutdown.changed() => {
        return Ok(());
    }
}

一个分支完成后,其他未完成 Future 会被 drop。Future 的 drop 通常表示取消等待,但取消后的业务状态是否一致,取决于具体操作。

适合直接取消的例子:

  • 等待 channel 下一条消息;
  • 等待 timer;
  • 等待 socket 下一次可读,且准备关闭整个连接。

需要额外设计的例子:

  • 多步写协议帧;
  • 已提交但未确认的远程请求;
  • 跨多个 await 的状态迁移;
  • 需要回滚或幂等保证的事务。

8.2 spawnselect! 的角色不同

问题tokio::spawntokio::select!
是否创建新 task
生命周期可独立结束与当前 task 相同
局部借用通常受 'static 限制可正常借用局部变量
结束观察JoinHandle / JoinSet获胜分支的结果
常见用途独立连接、后台管理任务超时、关闭、竞速、多输入状态机

8.3 在循环中继续同一个 Future

rust
let refresh = refresh_snapshot();
tokio::pin!(refresh);

loop {
    tokio::select! {
        result = &mut refresh => {
            println!("snapshot ready: {result:?}");
            break;
        }
        Some(command) = control_rx.recv() => {
            handle_control(command);
        }
    }
}

Future 在循环外创建,并通过 &mut refresh 继续轮询。如果每轮重新调用 refresh_snapshot(),工作会反复从头开始。


九、Stream:连续异步事件的消费接口

Future 最终至多产出一个值;Stream 可以多次产值,直到返回 None

rust
use tokio_stream::StreamExt;

let events = event_receiver
    .filter_map(|result| result.ok())
    .filter(|event| event.key.starts_with("user:"))
    .take(10);

tokio::pin!(events);

while let Some(event) = events.next().await {
    println!("changed: {}", event.key);
}

适配器顺序决定业务语义:

rust
stream.filter(is_relevant).take(10) // 等到 10 个相关事件
stream.take(10).filter(is_relevant) // 只检查最先到达的 10 个事件

Stream 只定义消费方式,不替你决定底层队列策略:

channel语义慢消费者后果
mpsc每条消息交给一个消费者发送方在容量满时等待
broadcast每个订阅者都收到消息慢订阅者可能发生 lag
watch只关心最新状态中间变化会被覆盖
oneshot单个结果没有连续流

在 PocketKV 中,可以用 broadcast 发布 key 变更事件,再用 BroadcastStream 暴露订阅接口。调用方必须处理 Lagged,不能假装所有事件都永久保存。


十、优雅关闭:触发、下发、回收、截止

10.1 关闭不等于 runtime 被销毁

直接让 main 返回会终止 runtime,尚未完成的 task 可能被丢弃。更可靠的流程是:

text
收到终止信号
   ▼
停止 accept 新连接
   ▼
广播 shutdown
   ▼
连接完成当前安全边界并退出
   ▼
JoinSet 汇总结果
   ▼
超过 deadline 时 abort 剩余 task

10.2 watch 下发最新关闭状态

rust
use tokio::sync::watch;

let (shutdown_tx, _) = watch::channel(false);

每条连接通过 shutdown_tx.subscribe() 获得 Receiver。使用 watch 的原因是关闭属于状态,而不是必须逐条消费的事件;后来创建的订阅者也能看到最新值。

10.3 JoinSet 回收任务

rust
use tokio::task::JoinSet;

let mut tasks = JoinSet::new();

tasks.spawn(serve_connection(
    socket,
    Arc::clone(&store),
    shutdown_tx.subscribe(),
));

运行期间可以顺手回收已结束 task,避免结果一直积压:

rust
Some(result) = tasks.join_next(), if !tasks.is_empty() => {
    match result {
        Ok(Ok(())) => {}
        Ok(Err(error)) => eprintln!("connection error: {error}"),
        Err(join_error) => eprintln!("connection task panicked: {join_error}"),
    }
}

10.4 完整服务器监督骨架

rust
use std::sync::Arc;
use tokio::net::TcpListener;
use tokio::sync::{watch, Semaphore};
use tokio::task::JoinSet;
use tokio::time::Duration;

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let listener = TcpListener::bind("127.0.0.1:6380").await?;
    let store = Arc::new(Store::default());
    let permits = Arc::new(Semaphore::new(512));
    let (shutdown_tx, _) = watch::channel(false);
    let mut tasks = JoinSet::new();

    loop {
        tokio::select! {
            accepted = listener.accept() => {
                let (socket, peer) = accepted?;
                let permit = Arc::clone(&permits).acquire_owned().await?;
                let store = Arc::clone(&store);
                let shutdown = shutdown_tx.subscribe();

                tasks.spawn(async move {
                    let _permit = permit;
                    let result = serve_connection(socket, store, shutdown).await;
                    if let Err(error) = &result {
                        eprintln!("peer {peer}: {error}");
                    }
                    result
                });
            }
            signal = tokio::signal::ctrl_c() => {
                signal?;
                println!("shutdown requested");
                break;
            }
            Some(result) = tasks.join_next(), if !tasks.is_empty() => {
                observe_connection_result(result);
            }
        }
    }

    let _ = shutdown_tx.send(true);

    let deadline = tokio::time::sleep(Duration::from_secs(10));
    tokio::pin!(deadline);

    while !tasks.is_empty() {
        tokio::select! {
            result = tasks.join_next() => {
                if let Some(result) = result {
                    observe_connection_result(result);
                }
            }
            _ = &mut deadline => {
                let remaining = tasks.len();
                eprintln!("shutdown deadline exceeded; aborting {remaining} tasks");
                tasks.abort_all();
                while tasks.join_next().await.is_some() {}
                break;
            }
        }
    }

    println!("PocketKV stopped");
    Ok(())
}

这里有一个值得审查的细节:示例在 accept 后等待 permit,短时间内可能持有已接收 socket。另一种写法是在 accept 前取得 permit,但若 accept 长期等待,permit 会被空占。可以使用 select! 同时协调 permit、accept 和 shutdown,或采用 try_acquire 后立即拒绝。架构选择应依据希望限制的是“已接收连接”还是“活跃处理任务”。

关闭等待没有把 JoinSet 移入单独的 drain Future,而是让 join_next() 与 deadline 在同一循环中竞速。这样截止时间到达时,主任务仍然拥有 JoinSet,可以调用 abort_all() 并继续回收取消结果。

完整关闭需要两个方向 watch 向下传播“请退出”,JoinSet 向上汇总“已退出”。只有通知没有确认,无法知道资源是否真的回收。


十一、同步世界如何复用异步客户端

11.1 简单阻塞包装

rust
use tokio::runtime::Runtime;

pub struct BlockingClient {
    runtime: Runtime,
    inner: AsyncClient,
}

impl BlockingClient {
    pub fn connect(address: &str) -> Result<Self, ClientError> {
        let runtime = tokio::runtime::Builder::new_current_thread()
            .enable_all()
            .build()?;
        let inner = runtime.block_on(AsyncClient::connect(address))?;

        Ok(Self { runtime, inner })
    }

    pub fn get(&mut self, key: &str) -> Result<Option<Bytes>, ClientError> {
        self.runtime.block_on(self.inner.get(key))
    }
}

这种包装适合命令行工具或调用频率不高的同步接口。current_thread runtime 只在 block_on 期间驱动任务;方法返回后,后台 spawned task 不会继续运行,直到下一次进入 runtime。

11.2 独立 runtime 线程

GUI、插件宿主或长期运行的同步服务可能需要异步系统持续工作。可以启动专用线程,让它拥有 runtime 和客户端,外部通过有界 channel 发请求:

text
同步调用线程 ──blocking_send──▶ runtime 线程
同步调用线程 ◀──blocking_recv── runtime 线程

这与前面的 manager task 是同一思想:资源只有一个所有者,跨边界使用消息协议。

11.3 异步代码调用阻塞工作

方向相反时,使用 spawn_blocking

rust
let decoded = tokio::task::spawn_blocking(move || decode_large_file(input))
    .await??;

不要在 async task 中直接调用 std::thread::sleep、长时间同步文件 I/O 或重 CPU 函数,否则会占住 runtime Worker,使同线程上的其他 task 无法推进。

11.4 避免嵌套 block_on

库代码不应在不知调用环境的情况下随意创建 runtime 并 block_on。优先公开 async API,把同步包装放到明确的边界模块,并说明线程、阻塞和关闭行为。


十二、错误类型要支持上层决策

rust
#[derive(Debug, thiserror::Error)]
pub enum ServerError {
    #[error("I/O error: {0}")]
    Io(#[from] std::io::Error),

    #[error("protocol error: {0}")]
    Protocol(#[from] ProtocolError),

    #[error("idle timeout")]
    IdleTimeout,

    #[error("server is shutting down")]
    ShuttingDown,

    #[error("server is overloaded")]
    Overloaded,
}

分层错误让调用者选择不同动作:

  • 协议错误:向客户端返回 Error Frame;
  • 单连接 I/O 错误:记录后关闭连接;
  • Overloaded:拒绝或重试;
  • ShuttingDown:停止启动新操作;
  • manager task 结束:可能需要触发服务级降级或关闭。

把所有失败压成字符串会丢掉可执行的语义,把可预期错误变成 panic 则会放大故障范围。


十三、测试策略:覆盖切片、取消和容量

13.1 Frame parser 表驱动测试

至少覆盖:

  • 完整 Simple/Bulk/Array;
  • 一个 Frame 分成多个切片;
  • 两个 Frame 连续到达;
  • 空缓存 EOF;
  • 半帧 EOF;
  • 长度字段溢出;
  • 超长 Bulk;
  • 递归过深;
  • 随机字节不 panic、不无限循环。

13.2 使用内存双端流测试 Connection

Tokio 提供 tokio::io::duplex,可在不监听真实端口的情况下模拟分片读取与写入:

rust
#[tokio::test]
async fn reads_frame_split_across_writes() {
    let (mut client, server) = tokio::io::duplex(64);

    let writer = tokio::spawn(async move {
        use tokio::io::AsyncWriteExt;
        client.write_all(b"$5\r\nhe").await.unwrap();
        client.write_all(b"llo\r\n").await.unwrap();
    });

    let mut connection = Connection::from_stream(server);
    let frame = connection.read_frame().await.unwrap().unwrap();
    assert_eq!(frame, Frame::Bulk(Bytes::from_static(b"hello")));

    writer.await.unwrap();
}

13.3 时间测试使用虚拟时钟

启用 Tokio 测试时间控制后,可以暂停并推进 timer,避免测试真的等待 30 秒。超时测试应验证:Future 被取消后连接和状态仍处于预期状态。

13.4 Loom 与并发状态

如果实现复杂的原子或自定义同步结构,可以用 Loom 探索线程交错。普通 Mutex<HashMap> 不必为了“高级”而改成无锁结构,除非 profile 已证明它是热点。


十四、设计取舍复盘

问题本章基线何时升级可能方向
连接并发每连接一个 task + Semaphore需要租户隔离或分级限流多级 permit、准入控制
简单状态Arc<Store> + 同步 Mutex临界区竞争明显分片、专门 manager
严格顺序资源manager + 有界 mpsc单 manager 成为瓶颈分区 manager、连接池
单次返回oneshot一个请求产生多次值mpsc / Stream
状态广播watch需要保留每次事件broadcast / 持久日志
task 管理JoinSet任务有父子层级与重启策略supervisor、结构化并发
协议缓冲BytesMut + 上限极端吞吐或零拷贝需求基准后调整、批量处理
同步桥接Runtime::block_on后台任务必须持续运行专用 runtime 线程

Tokio 提供调度器、I/O reactor、timer 和同步原语,但它不会替你确定业务容量、过载策略、取消一致性与关闭期限。


十五、常见失误

15.1 忘记 .await

调用 async 函数后只得到 Future,网络操作还没有执行。不要忽略 must_use 警告。

15.2 spawn 后完全不观察结果

丢弃 JoinHandle 会让业务错误和 panic 难以被监督。普通连接至少记录错误,关键后台任务应由 JoinSet 或 supervisor 管理。

15.3 跨 .await 持有同步锁

症状包括 Future 不满足 Send、吞吐下降和其他 task 长时间阻塞。把状态操作封装成同步方法,缩短 guard 生命周期。

15.4 把异步 Mutex 当作默认答案

tokio::sync::Mutex 允许跨 .await 持锁,但临界区仍然串行。只有资源确实需要在锁保护下等待异步操作时才考虑它,并审查取消后状态。

15.5 使用无界 channel

短期压力被积累成长期内存和尾延迟。使用有界容量,并为“满”建立指标和策略。

15.6 假设一次 read 就有一个 Frame

TCP 没有消息边界。必须保留缓冲并区分完整、不完整、非法和半帧 EOF。

15.7 忽略 select! 的取消后果

获胜分支执行时,其余 Future 被 drop。不可安全取消的操作应拆成独立 task、状态机或幂等步骤。

15.8 在循环中反复创建长任务 Future

每轮都重新开始会让任务永远完不成。需要继续同一 Future 时,在循环外创建并 pin。

15.9 在 async task 里调用阻塞函数

同步睡眠、重 CPU 和慢文件操作会堵住 runtime Worker。改用异步 API 或 spawn_blocking

15.10 把“发送关闭”当成“已经关闭”

广播只能表达意图。仍需等待 task、处理错误,并设置截止时间。


十六、练习

练习 1:准入策略对比

分别实现“等待 permit”“立即拒绝”“等待 50ms 后拒绝”,在固定并发下比较 P50/P99、成功率和内存。

练习 2:完成协议层

实现全部 Frame variant 的检查、解析和写入。为数字溢出、递归深度和非法 CRLF 增加错误类型。

练习 3:命令解析矩阵

用表驱动测试覆盖命令名大小写、缺少参数、多余参数、非 Bulk key、空 key 与未知命令。任何输入都不得 panic。

练习 4:带过期时间的 SET

增加 SET key value PX milliseconds。比较以下实现:

  • 每个 key 一个 timer task;
  • 最小堆 + 单个维护 task;
  • 惰性删除。

说明内存、精度和关闭行为。

练习 5:Manager 批处理

让 manager 每次接收第一条命令后,用 try_recv 收集一小批请求再统一处理。通过 benchmark 判断吞吐提升是否值得增加尾延迟。

练习 6:订阅变更事件

提供 subscribe(prefix) -> impl Stream<Item = Result<Event, SubscribeError>>。定义慢订阅者发生 lag 时是报错、跳到最新还是断开。

练习 7:取消安全审计

列出连接循环中每个 .await。对每一点说明 Future 被 drop 后:

  • socket 是否可继续使用;
  • 是否可能写出半帧;
  • Store 是否已变更;
  • 客户端是否可能重试并造成重复操作。

练习 8:两阶段关闭

收到 Ctrl+C 后停止 accept,允许连接完成当前命令但不读取下一条命令;10 秒后 abort 剩余 task,并输出未正常结束数量。

练习 9:同步客户端边界

分别实现 BlockingClient 与专用 runtime 线程版本,记录它们对后台 task、线程安全、关闭和嵌套 runtime 的约束。


十七、总结

  1. async 函数创建 Future,spawn 才创建独立 task;.await 本身不等于并发。
  2. Tokio task 很轻量,但连接、缓冲、channel 和状态仍需要明确上限。
  3. Semaphore 可以让连接准入与 task 生命周期绑定。
  4. TCP 是字节流,Connection 必须保存未消费数据,并处理完整帧、半帧和连续帧。
  5. Frame 表达协议语法,Command 表达合法业务输入,Store 表达状态责任。
  6. 短同步临界区可以使用 std::sync::Mutex,但 guard 不能跨 .await
  7. 单一 manager task 适合严格顺序或独占资源,通过有界 mpsc 接收命令、oneshot 返回单次结果。
  8. channel 容量是背压合同,不只是调优参数。
  9. select! 选择一个分支并取消其他 Future,必须审查 cancellation safety。
  10. Stream 适合连续事件,但底层消息保留与慢消费者策略由 channel 决定。
  11. 优雅关闭包括停止入口、向下通知、向上回收和截止时间。
  12. JoinSet 让 spawned task 的错误与 panic 不再成为不可见背景状态。
  13. 同步桥接应位于明确边界;异步代码中的阻塞工作应交给 spawn_blocking
  14. Tokio 解决运行时基础设施,不会自动解决容量、过载、状态不变量和业务幂等。

章节导航

  • 上一章:第13章:多线程 Web 服务器
  • 下一章:第15章:可靠高性能 Rust

  • Rust 深度教程
  • 第9章:并发编程
  • 第10章:异步编程
  • 第12章:Cargo 工程化
  • 第13章:多线程 Web 服务器
  • 第15章:可靠高性能 Rust

本章目录
导读:异步不是把线程换成 .await一、项目依赖与 feature 预算二、把连接接入变成受控并发三、协议分层:字节、Frame、Command 各司其职四、命令类型把合法输入收窄五、共享状态方案一:短临界区的同步 Mutex六、共享状态方案二:单一所有者与消息协议七、连接循环:将 I/O、解析、执行串起来八、select! 是取消边界,不只是语法糖九、Stream:连续异步事件的消费接口十、优雅关闭:触发、下发、回收、截止十一、同步世界如何复用异步客户端十二、错误类型要支持上层决策十三、测试策略:覆盖切片、取消和容量十四、设计取舍复盘十五、常见失误十六、练习十七、总结章节导航Related Documents
苏ICP备2025204887号-2