Learn Rust Coding · 课时

Crossbeam 通道

高级通道

第 4 / 4 课13 个步骤

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 反馈 — 无需本地设置。

此课程中的所有课时

  1. mpsc 通道
  2. 使用 Arc/Mutex 共享状态
  3. 作用域线程
  4. Crossbeam 通道
← 返回 Learn Rust Coding