两种并发模型
Rust 提供两种并发模型,各有适用场景:
| 模型 | 适用场景 | Rust 工具 |
|---|---|---|
| OS 线程 | CPU 密集型、并行计算 | std::thread + Rayon |
| 异步任务 | I/O 密集型、高并发网络 | async/await + Tokio |
OS 线程与消息传递
use std::thread;
use std::sync::mpsc; // multi-producer, single-consumer
let (tx, rx) = mpsc::channel();
let tx2 = tx.clone(); // 多个发送端
thread::spawn(move || {
tx.send(String::from("来自线程 1")).unwrap();
});
thread::spawn(move || {
tx2.send(String::from("来自线程 2")).unwrap();
});
// 接收两条消息
for msg in rx {
println!("收到: {}", msg);
}move 关键字让闭包获取变量的所有权,而非借用。新线程的生命周期不确定,编译器拒绝借用——这就是 Rust 用所有权防止数据竞争的方式。
智能指针:堆上的数据
Box<T>:堆分配
// 编译时大小未知的递归类型必须用 Box
enum List {
Cons(i32, Box<List>),
Nil,
}
let list = List::Cons(1,
Box::new(List::Cons(2,
Box::new(List::Nil))));Arc<T> + Mutex<T>:线程安全的共享状态
对比:
JS
JavaScript// JS 单线程,共享状态天然安全
let counter = { value: 0 };
// 直接修改,无需同步原语Rs
Rustuse std::sync::{Arc, Mutex};
use std::thread;
let counter = Arc::new(Mutex::new(0));
let mut handles = vec![];
for _ in 0..10 {
let counter = Arc::clone(&counter);
let handle = thread::spawn(move || {
let mut num = counter.lock().unwrap();
*num += 1;
});
handles.push(handle);
}
for h in handles { h.join().unwrap(); }
println!("结果: {}", *counter.lock().unwrap()); // 10| 类型 | 用途 | 线程安全 |
|---|---|---|
Rc<T> | 单线程引用计数 | ❌ |
Arc<T> | 多线程引用计数 | ✅ |
RefCell<T> | 单线程内部可变性 | ❌ |
Mutex<T> | 多线程互斥锁 | ✅ |
RwLock<T> | 多线程读写锁 | ✅ |
async/await:异步编程
对比:
JS
JavaScriptasync function fetchUser(id) {
const res = await fetch(`/users/${id}`);
return res.json();
}
// 并发执行
const [u1, u2] = await Promise.all([
fetchUser(1), fetchUser(2)
]);Rs
Rustasync fn fetch_user(id: u64) -> User {
reqwest::get(format!("/users/{id}"))
.await.unwrap()
.json::<User>()
.await.unwrap()
}
// 并发执行
let (u1, u2) = tokio::join!(
fetch_user(1),
fetch_user(2)
);Rust 的 async fn 返回 Future,Future 是惰性的——必须被运行时(如 Tokio)驱动才会执行。.await 暂停当前任务,让运行时执行其他任务,不阻塞线程。
Tokio:异步运行时
use tokio::time::{sleep, Duration};
#[tokio::main]
async fn main() {
// 生成独立任务(类比 Promise,但更轻量)
let handle = tokio::spawn(async {
sleep(Duration::from_millis(100)).await;
42
});
let result = handle.await.unwrap();
println!("任务结果: {}", result);
}tokio::select!:等待多个 Future
tokio::select! {
result = fetch_from_api() => {
println!("API 响应: {:?}", result);
}
_ = timeout(Duration::from_secs(5)) => {
println!("超时!");
}
}Rayon:数据并行
use rayon::prelude::*;
let numbers: Vec<i32> = (0..1_000_000).collect();
// 自动利用所有 CPU 核心
let sum: i32 = numbers.par_iter()
.map(|&x| x * x)
.sum();将 .iter() 换成 .par_iter() 就能启用并行——Rayon 自动处理工作分配和线程池。适合 CPU 密集型批处理,I/O 密集型用 Tokio。
并发爬虫构建演示
从单线程爬虫到 Tokio 异步并发爬虫:展示 async/await 如何用少量线程处理大量并发 HTTP 请求,对比性能差异
视频即将上线
实战项目
并发爬虫工具
初级
构建一个异步并发网页爬虫:给定 URL 列表,用 Tokio + reqwest 并发抓取,Arc<Mutex> 汇总结果,最终输出各页面的标题、状态码和响应时间统计。
tokioasync/awaitArc/Mutexreqwest并发控制