Crossbeamチャネル
高度なチャネル
「Crossbeamチャネル」はCoddyKit上の無料Learn Rust Codingレッスンです。 これはレッスン4/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはLearn Rust Coding学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 Learn Rust Codingコースには全4レッスンが含まれています。
std mpscの先へ
標準のmpscチャネルは単一コンシューマーです。Receiverは1つしかありません。crossbeam-channelクレートは、より多くの機能と、より優れたパフォーマンスを提供することが多いマルチプロデューサー・マルチコンシューマー(mpmc)チャネルを提供します。
主な特徴は次のとおりです。
- senderだけでなく、clone可能な
Receiver。 - 多数のチャネルを待機できる強力な
select!マクロ。 - 容量制限付き、無制限、tickなどの特殊なチャネル。
依存関係の追加
Crossbeamは外部クレートなので、Cargo.tomlに追加します。
[dependencies]crossbeam-channel = "0.5"
続いて、その関数をimportします。Cargoと外部クレートが必要なため、ここで示すスニペットは単体で実行するものではなく、APIの使い方を示すものです。
// 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());
}unboundedとbounded
Crossbeamには主に2つのコンストラクターがあります。
unbounded():必要に応じて拡張され、sendはブロックしません。bounded(cap):固定サイズのバッファーを持ち、満杯になるとsendがブロックしてバックプレッシャーを発生させます。
bounded(0)チャネルはランデブーチャネルで、sendとrecvが直接受け渡しを行います。
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をcloneできることです。複数のワーカースレッドが同じチャネルから取り出せ、各メッセージはそのうちの1つだけに渡されます。これは、ワークスティーリング型スレッドプールの基盤となります。
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!を使うと、1つのスレッドで複数のチャネル操作を待機し、最初に準備できた操作に応じて処理できます。チャネルイベントに対する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)関数は、一定間隔でメッセージを届けるReceiverを返します。これを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はパイプラインの構築に適しています。ステージ1が生成し、ステージ2が変換し、ステージ3が消費します。各ステージはチャネルで接続された独自のスレッドで実行され、受信側をクローンすれば、任意のステージを複数のワーカーに拡張できます。
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を使用します。
- 1つのキューを複数のコンシューマーで共有する場合(ワーカープール)。
- タイムアウト付きで複数のチャネルを
select!する場合。 - 選択処理に組み込まれた
tickによる定期タイマーを使用する場合。
1つのプロデューサーと1つのコンシューマーによる単純なパイプであれば、標準のmpscで十分であり、依存関係も必要ありません。
チャネルはドロップによって閉じられる
stdと同様に、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を使用します。
よくある質問
「Crossbeamチャネル」レッスンは無料ですか?
はい。「Crossbeamチャネル」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、Learn Rust Codingコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 Learn Rust Codingコースには全4レッスンが含まれています。
「Crossbeamチャネル」で何を学びますか?
高度なチャネル ブラウザで直接実行するハンズオンコードでLearn Rust Codingを演習し、24時間対応の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チャネル