ARTICLE DETAIL

资讯详情

深耕网站视觉设计与运营推广的一线实战洞察。

Rust 异步取消(Cancellation)实战:tokio `select!` 中的数据丢失陷阱与取消安全设计

Rust 异步取消(Cancellation)实战:tokio `select!` 中的数据丢失陷阱与取消安全设计 Rust 异步取消Cancellation实战tokioselect!中的数据丢失陷阱与取消安全设计【免费下载链接】comprehensive-rustThis is the Rust course used by the Android team at Google. It provides you the material to quickly teach Rust.项目地址: https://gitcode.com/GitHub_Trending/co/comprehensive-rust本文基于 Google Android 团队维护的开源 Rust 课程 comprehensive-rust 的并发模块「异步陷阱」章节src/concurrency/async-pitfalls/cancellation.md展开。课程由 Android 团队用于快速讲授 Rust本文聚焦该章节核心主题当 Future 在任意await点被 drop 时即发生取消cancellation。你将通过一个可运行的 tokio 示例理解取消的触发机制、为什么编译器无法帮你发现取消安全问题以及如何把易丢失的局部状态移入结构体来编写取消安全cancellation-safe的异步代码并掌握Interval::tick、AsyncReadExt::read、AsyncBufReadExt::read_line等常用 API 的取消安全性差异。取消是什么drop 一个 Future 意味着它永远不会再被 poll在 Rust 的 async/await 模型中一个async fn或async块在被调用时会创建一个实现了Futuretrait 的类型其中捕获了该函数的全部局部变量关于这一点课程在 Pin 章节中做了详细解释——Future 内部可能持有指向自身其他局部变量的自引用指针因此它必须通过Pin才能被安全地 poll。取消cancellation的定义十分朴素drop 一个 Future 意味着它永远不会再被 poll 了。由于异步代码在每次.await处都可能让出执行权因此取消可能发生在任何一个 await 点。例如使用tokio::select!时竞争失败的其它分支对应的 Future 会被直接 drop持有JoinHandle的任务被 drop或运行时被关闭时tokio::spawn产生的任务会被取消未来在循环中重新创建/替换时旧的 Future 被 drop。课程的 tokio 运行时章节就给出了一个活生生的取消例子tokio::spawn(count_to(10))后main任务循环 5 次即结束count_to永远到不了 10——因为 main 结束后运行时关闭spawn 出去的任务连同它的 Future 一起被取消这就是异步取消。因此设计异步系统时必须格外小心即使 Future 被取消系统也必须保持正确——既不能死锁也不能丢数据。这正是本节课程要训练的思维。一个会丢数据的完整示例慢速拷贝 行读取器课程给出的示例由两个部分构成一个逐字节写入的慢速拷贝任务和一个通过DuplexStream逐字节累积字符串的行读取器二者通过tokio::select!与一个 60ms 的定时器 tick 竞争。完整代码如下课程中标记为compile_fail刻意留作课堂讨论素材核心教学点是运行时的数据丢失现象use std::io; use std::time::Duration; use tokio::io::{AsyncReadExt, AsyncWriteExt, DuplexStream}; struct LinesReader { stream: DuplexStream, } impl LinesReader { fn new(stream: DuplexStream) - Self { Self { stream } } async fn next(mut self) - io::ResultOptionString { let mut bytes Vec::new(); let mut buf [0]; while self.stream.read(mut buf[..]).await? ! 0 { bytes.push(buf[0]); if buf[0] b\n { break; } } if bytes.is_empty() { return Ok(None); } let s String::from_utf8(bytes) .map_err(|_| io::Error::new(io::ErrorKind::InvalidData, not UTF-8))?; Ok(Some(s)) } } async fn slow_copy(source: String, mut dest: DuplexStream) - io::Result() { for b in source.bytes() { dest.write_u8(b).await?; tokio::time::sleep(Duration::from_millis(10)).await } Ok(()) } #[tokio::main] async fn main() - io::Result() { let (client, server) tokio::io::duplex(5); let handle tokio::spawn(slow_copy(hi\nthere\n.to_owned(), client)); let mut lines LinesReader::new(server); let mut interval tokio::time::interval(Duration::from_millis(60)); loop { tokio::select! { _ interval.tick() println!(tick!), line lines.next() if let Some(l) line? { print!({}, l) } else { break }, } } handle.await.unwrap()?; Ok(()) }逐段拆解这个程序的工作原理tokio::io::duplex(5)创建一对互相连接的DuplexStream内部缓冲区容量为 5 字节。client端交给slow_copy写入server端由LinesReader读取。slow_copy把字符串hi\nthere\n9 个字节逐字节写入client每写一个字节就sleep(10ms)。由于缓冲区只有 5 字节写入会因为背压而阻塞等待读取方消化数据从而保证两个任务交替推进、竞争真实发生。LinesReader::next每次调用都新建局部的bytes: Vecu8和buf: [u8; 1]逐字节读入遇到\n才返回一行字符串读到流末尾且无数据时返回Ok(None)表示结束。tokio::select!每一轮循环在interval.tick()每 60ms 一次和lines.next()之间竞争谁先就绪就执行谁的分支未就绪的另一个 Future 被 drop。关于select!的详细语义可参考课程 Select 章节。运行这个程序输出会表现为打印若干tick!且最终打印出来的行是残缺的——例如可能只有hi或ther的一部分。字符串的很大一部分被悄悄丢掉了。为什么数据会丢局部状态随 Future 一起被 drop秘密藏在LinesReader::next的实现里async fn next(mut self) - io::ResultOptionString { let mut bytes Vec::new(); // 局部变量累积已读字节 let mut buf [0]; // 局部变量当前读入的 1 字节 while self.stream.read(mut buf[..]).await? ! 0 { bytes.push(buf[0]); if buf[0] b\n { break; } } // ... }bytes和buf都是next的局部变量它们被存储在next生成的 Future 内部。当select!中interval.tick()先就绪时未完成的next()Future连同它内部的bytes、buf以及从stream读取的部分数据被整体drop下次循环再次调用next()时一切从零开始新建空的bytes重新从流中读取后续字节。于是在两次next()调用之间被吃掉的字节就永久丢失了。这正是课程强调的核心结论只要tick()分支先完成next()和它的buf就被 drop示例就丢失字符串的一部分。编译器不会帮你取消安全需要人工审查课程在补充材料中特别强调了两个关键认知编译器不会帮助你检查取消安全性。与借用检查器在编译期严格守护内存安全不同没有任何静态分析会告诉你这个 Future 在 await 点被取消后你累积的中间状态会丢。你需要阅读 API 文档并仔细审视自己的async fn内部持有哪些状态自行判断取消后这些状态是否仍然完整。取消与panic、?错误传播有本质区别。panic和?属于错误处理范畴是异常路径而取消是正常控制流的一部分——Future 被 drop 是合法、频繁、默认发生的例如select!每轮都要 drop 失败的分支。因此不能把取消当成错误来捕获而必须把可取消当作代码的常态假设来设计。从更广的角度看这也呼应了课程异步陷阱章节src/concurrency/async-pitfalls.md的立意async/await 提供了便捷高效的并发抽象但也带来了若干脚枪footgun取消安全正是其中最隐蔽的一个。修复方案把可变状态移入结构体要让LinesReader变得取消安全思路是把每次调用之间需要持久保留的中间状态从next的局部变量提升为结构体字段即mut self的一部分这样 Future 被 drop 时状态并不会丢下次调用可以从上次中断处继续。课程给出的改造如下struct LinesReader { stream: DuplexStream, bytes: Vecu8, // 累积的已读字节从局部变量提升为字段 buf: [u8; 1], // 单字节缓冲区同样提升为字段 } impl LinesReader { fn new(stream: DuplexStream) - Self { Self { stream, bytes: Vec::new(), buf: [0] } } async fn next(mut self) - io::ResultOptionString { // prefix buf and bytes with self. // ... let raw std::mem::take(mut self.bytes); let s String::from_utf8(raw) .map_err(|_| io::Error::new(io::ErrorKind::InvalidData, not UTF-8))?; // ... } }改造要点说明bytes与buf成为LinesReader的字段后它们的所有者是mut self而next只借用mut self。即使next生成的 Future 被 drop这些字段也依然存活在LinesReader实例中读取进度被完整保留。每次成功累积到一行后用std::mem::take(mut self.bytes)把累积内容取走并替换为空Vec既避免克隆开销又为下一行的累积腾出干净起点。由于stream本身也是self的字段已读走的字节不会再被读到下次next()继续从流的当前位置读取从而保证逐行数据不重不漏。这个模式可以推广为一条通用的取消安全设计准则凡是跨 await 必须持久的累积状态缓冲区、游标、计数、中间结果都应存放在 Future 之外的持有者结构体字段、ArcMutex..、任务内部循环外的变量等中而不是async fn的局部变量里。常用 API 的取消安全性对照课程的补充材料直接点明了三个 tokio 常用 API 的取消安全性这是实战中判断能否在select!里放心使用的关键依据API取消安全性原因Interval::tick✅ 安全它内部跟踪了一个 tick 是否已投递delivered取消后再次 poll 不会重复或丢失周期AsyncReadExt::read✅ 安全它的语义是要么返回、要么完全不读数据取消不会吞掉半截读取AsyncBufReadExt::read_line❌ 不安全与本文示例的LinesReader::next类似它会先把数据读进内部缓冲再返回取消会丢失已读内容官方文档对此有详细说明并给出了替代方案把这三点记成直觉读一行/累积式的 API 天然容易在取消时丢状态使用时要么查阅文档确认其取消安全性要么自行承担缓冲职责如上面的结构体方案而单次、原子的 API如read读满给定缓冲通常安全。另外需要留意取消安全的tick也意味着它在select!循环中每 60ms 触发一次的节奏是可靠的而如果你把一个一次性的sleepFuture 放进select!循环如 Pin 章节中的 actor 示例一旦它Poll::Ready就会每轮都就绪——课程用Box::pin 每次到期后重新赋值timeout_fut的方式解决这正是取消/重复 poll 语义在循环场景下的另一种坑。综合实践建议综合本节课程编写生产级异步代码时可以遵循以下检查清单默认假设任何.await都可能被取消尤其是select!分支、tokio::spawn的任务、以及被超时包装的 Future。审视跨 await 的局部状态async fn里的Vec、计数器、游标、半成品缓冲取消即丢需要持久的一律提升到mut self、外部容器或任务级循环之外。优先使用文档标注为取消安全的 API如Interval::tick、AsyncReadExt::read避免read_line这类累积式 API 带来的隐性丢数据。把取消当作正常控制流设计而不是错误处理不要试图捕获取消而要保证取消后系统状态一致、无死锁、无数据丢失。结合Pin与select!的语义理解运行时机Future 必须通过Pin才能 poll见 Pin 章节select!每轮只保留一个赢家、drop 其余分支见 Select 章节。想同时等待多个 Future 且不允许取消丢失的则应考虑join!/join_all的语义见 Join 章节。本小节在课程中用于约 18 分钟的课堂讲解frontmatter 中minutes: 18是异步陷阱四连击Cancellation、Pin、Blocking the executor、Async traits的第一课读者可结合 src/concurrency/async-pitfalls/pin.md、src/concurrency/async-pitfalls/blocking-executor.md 继续深入或在 src/SUMMARY.md 中定位整门课程的学习路径。【免费下载链接】comprehensive-rustThis is the Rust course used by the Android team at Google. It provides you the material to quickly teach Rust.项目地址: https://gitcode.com/GitHub_Trending/co/comprehensive-rust创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表