Crossbeam 通道
高级通道
Crossbeam 通道 是 CoddyKit 上的免费 Learn Rust Coding 课时。 这是第 4 节课,共 4 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 AI 导师进行实践。 这是 Learn Rust Coding 学习路径的一部分,你的进度在网页和 CoddyKit 应用中同步。 Learn Rust Coding 课程共包含 4 节课。
超越标准库 mpsc
标准库中的 mpsc 通道是单消费者通道:只有一个接收器。crossbeam-channel 代码包提供多生产者、多消费者(mpmc)通道,功能更多,性能通常也更好。
主要特点:
- 可以克隆的
Receiver,而不仅仅是发送器。 - 功能强大的
select!宏,可等待多个通道。 - 有界、无界通道,以及滴答通道等特殊通道。
添加依赖
Crossbeam 是外部代码包,因此请将其添加到 Cargo.toml:
[dependencies]crossbeam-channel = "0.5"
然后导入其中的函数。由于这里需要 Cargo 和外部代码包,以下代码片段用于说明应用程序接口,而不是作为独立程序运行。
// Cargo.toml
// [dependencies]
// crossbeam-channel = "0.5"
use crossbeam_channel::unbounded;
fn main() {
let (s, r) = unbounded();
s.send("hello").unwrap();
println!("{}", r.recv().unwrap());
}无界通道与有界通道
Crossbeam 提供两个主要的构造函数:
unbounded()会按需增长;send永远不会阻塞。bounded(cap)具有固定大小的缓冲区;缓冲区已满时,send会阻塞,从而提供反压。
bounded(0) 通道是一个会合通道,发送和接收会直接交接。
use crossbeam_channel::bounded;
use std::thread;
fn main() {
let (s, r) = bounded(2);
thread::spawn(move || {
for i in 1..=3 {
s.send(i).unwrap();
}
});
while let Ok(v) = r.recv() {
println!("got {}", v);
}
}多个消费者
其关键优势在于:您可以克隆 Receiver。多个工作线程可以从同一个通道中获取任务,每条消息只会交给其中一个线程。这是工作窃取线程池的基础。
use crossbeam_channel::unbounded;
use std::thread;
fn main() {
let (s, r) = unbounded();
let mut workers = vec![];
for id in 0..3 {
let rx = r.clone();
workers.push(thread::spawn(move || {
while let Ok(job) = rx.recv() {
println!("worker {} got job {}", id, job);
}
}));
}
for job in 0..6 { s.send(job).unwrap(); }
drop(s);
for w in workers { w.join().unwrap(); }
}select! 宏
select! 允许一个线程等待多个通道操作,并处理最先就绪的操作。它类似于针对通道事件的 match,非常适合组合来自多个来源的输入。
use crossbeam_channel::{unbounded, select};
use std::thread;
fn main() {
let (s1, r1) = unbounded();
let (s2, r2) = unbounded();
thread::spawn(move || { s1.send("from one").unwrap(); });
thread::spawn(move || { s2.send("from two").unwrap(); });
for _ in 0..2 {
select! {
recv(r1) -> msg => println!("r1: {}", msg.unwrap()),
recv(r2) -> msg => println!("r2: {}", msg.unwrap()),
}
}
}使用 select! 设置超时
select! 支持 default 分支和 recv(after(duration)) 计时器分支,因此您可以在超时后放弃等待,而不是永远阻塞。这对于保持系统响应非常重要。
use crossbeam_channel::{unbounded, select, after};
use std::time::Duration;
fn main() {
let (_s, r) = unbounded::<i32>();
select! {
recv(r) -> msg => println!("received {:?}", msg),
recv(after(Duration::from_millis(100))) -> _ => {
println!("timed out waiting for a message");
}
}
}try_send 与 try_recv
与标准库一样,Crossbeam 也提供非阻塞版本。若有界通道已满,try_send 会立即失败;若当前没有可用消息,try_recv 会立即失败。二者都会返回描述性错误,您可以使用 match 对其进行处理。
use crossbeam_channel::{bounded, TrySendError};
fn main() {
let (s, _r) = bounded(1);
s.send(1).unwrap();
match s.try_send(2) {
Ok(()) => println!("sent"),
Err(TrySendError::Full(v)) => println!("full, kept {}", v),
Err(TrySendError::Disconnected(v)) => println!("closed, kept {}", v),
}
}tick 与周期性工作
tick(duration) 函数会返回一个接收器,该接收器按照固定时间间隔传递消息。将它与 select! 结合使用,即可让周期性任务与其他通道并行运行,例如心跳或轮询循环。
use crossbeam_channel::{tick, select};
use std::time::Duration;
fn main() {
let ticker = tick(Duration::from_millis(50));
let mut count = 0;
while count < 3 {
select! {
recv(ticker) -> _ => {
count += 1;
println!("tick {}", count);
}
}
}
}流水线中的各个阶段
Crossbeam 非常适合构建流水线:第一阶段生成数据,第二阶段进行转换,第三阶段消费数据。每个阶段都运行在独立线程中,并通过通道连接;克隆接收端后,还可以将任意阶段扩展为多个工作线程。
use crossbeam_channel::unbounded;
use std::thread;
fn main() {
let (in_s, in_r) = unbounded();
let (out_s, out_r) = unbounded();
thread::spawn(move || {
for n in in_r { out_s.send(n * n).unwrap(); }
});
for n in 1..=4 { in_s.send(n).unwrap(); }
drop(in_s);
for sq in out_r { println!("square: {}", sq); }
}Crossbeam 优于 std 的场景
在以下情况下,您可以选择 crossbeam-channel:
- 多个消费者共享一个队列(工作线程池)。
- 使用
select!同时等待多个通道,并设置超时。 - 通过集成到选择机制中的 周期计时器
tick。
如果只是需要一个简单的单生产者单消费者管道,标准库中的 mpsc 就足够了,而且无需添加依赖。
通过丢弃关闭通道
与标准库一样,当接收端的所有发送者,或发送端的所有接收者都被丢弃时,crossbeam 通道就会关闭。最后一个发送者被丢弃后,对接收端的迭代也会结束。请始终使用 drop 或通过作用域管理发送者,以便工作线程循环能够正常退出。
use crossbeam_channel::unbounded;
use std::thread;
fn main() {
let (s, r) = unbounded();
let h = thread::spawn(move || {
let total: i32 = r.iter().sum();
println!("total: {}", total);
});
for i in 1..=5 { s.send(i).unwrap(); }
drop(s); // closes channel so the loop ends
h.join().unwrap();
}快速检查
测试您对 crossbeam 通道的理解。
总结
您学习了crossbeam 通道:
- 它们是 mpmc 通道:发送者和接收者都可以被克隆。
unbounded()和bounded(n)用于控制缓冲和背压。select!会等待多个通道,并可通过after设置超时、通过周期性的tick进行等待。try_send/try_recv都是非阻塞的。- 工作线程池和复杂流水线适合使用 crossbeam;简单管道则适合使用 std mpsc。
用 AI 导师学习 Rust — 免费
在浏览器中编写并运行真实代码,获得全天候 AI 导师的即时帮助,并在网页或应用中继续学习。
- 课程
- 39
- 课程
- 144
常见问题解答
「Crossbeam 通道」课时是免费的吗?
是的 — 「Crossbeam 通道」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 Learn Rust Coding 课程的其余内容,请升级到 CoddyKit PRO。 Learn Rust Coding 课程共包含 4 节课。
「Crossbeam 通道」这节课中我会学到什么?
高级通道 你通过在浏览器中直接运行的动手代码来练习 Learn Rust Coding,全天候 AI 导师会在你学习这节课的过程中回答你的问题。
学习 Learn Rust Coding 需要有经验吗?
无需任何先前经验。CoddyKit 上的 Learn Rust Coding 课程适合初学者到高级学习者,你可以从这里开始或从头开始,按照自己的节奏学习。 这是第 4 节课,共 4 节。
「Crossbeam 通道」课时需要多长时间?
大多数 CoddyKit 课程大约需要 5–10 分钟。每节课都很精短且互动,所以你能稳步进步,并在网页和应用中从离开的地方继续。
我能在这节 Learn Rust Coding 课中编写并运行代码吗?
能。每节 Learn Rust Coding 课都包含内置代码编辑器,你可以在浏览器中直接编写并运行真实代码,并获得即时 AI 反馈 — 无需本地设置。
此课程中的所有课时
- mpsc 通道
- 使用 Arc/Mutex 共享状态
- 作用域线程
- Crossbeam 通道