mpscチャネル
スレッド間で送信します
「mpscチャネル」はCoddyKit上の無料Learn Rust Codingレッスンです。 これはレッスン1/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはLearn Rust Coding学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 Learn Rust Codingコースには全4レッスンが含まれています。
チャネルとは何ですか
チャネルは、あるスレッドから別のスレッドへ値を送る一方向のパイプです。Rust の標準ライブラリには std::sync::mpsc があり、mpsc は multiple producer, single consumer(複数の送信者、単一の受信者)を意味します。
- Sender 側は値を送り込みます。
- Receiver 側は値を取り出します。
チャネルを使うと、メモリを直接共有する代わりにメッセージを渡すことでスレッド間通信ができ、多くのデータ競合を避けられます。
チャネルを作成する
mpsc::channel() を呼び出すと、(Sender, Receiver) のタプルを取得できます。ここでは、生成したスレッドからメインスレッドへ 1 つの値を送ります。
use std::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
tx.send(42).unwrap();
});
let received = rx.recv().unwrap();
println!("Got: {}", received);
}send と recv
tx.send(value) は Result を返します。失敗するのは、受信者が破棄された場合だけです。rx.recv() は値が届くまでブロックし、すべての送信者がなくなると Err を返します。
sendは値の所有権をチャネルへ移します。recvは反対側で所有権を取り出します。
use std::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let msg = String::from("hello from thread");
tx.send(msg).unwrap();
});
let text = rx.recv().unwrap();
println!("{}", text);
}所有権はチャネルを通って移動する
send は値を値渡しで受け取るため、送信後にその値を使うことはできません。このコンパイル時の規則により、別のスレッドが所有するようになったデータへの古い参照を、どのスレッドも保持しないことが保証されます。
以下では、send の後に msg を表示しようとするとコンパイルエラーになるため、1 回だけ使用しています。
use std::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
let data = vec![1, 2, 3];
tx.send(data).unwrap();
// data is moved; cannot use it here
});
let v = rx.recv().unwrap();
println!("Sum: {}", v.iter().sum::<i32>());
}Receiver を反復処理する
Receiver は IntoIterator を実装しています。Receiver をループすると、チャネルが閉じるまで(すべての送信者が破棄されるまで)各値が生成されます。メッセージのストリームを消費する慣用的な方法です。
use std::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
for i in 1..=3 {
tx.send(i).unwrap();
}
});
for received in rx {
println!("Received: {}", received);
}
}clone による複数の送信者
mpsc の mp は、複数の送信者を持てることを意味します。Sender をクローンし、それぞれのスレッドにコピーを渡します。受信者は、すべてのクローンが破棄されるまで値を収集します。
use std::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel();
let tx2 = tx.clone();
thread::spawn(move || { tx.send("from A").unwrap(); });
thread::spawn(move || { tx2.send("from B").unwrap(); });
for msg in rx {
println!("{}", msg);
}
}チャネルを閉じるときの動作
すべての送信者が破棄されると、受信側のループは自動的に終了します。Sender が 1 つでも生きていると、for msg in rx はさらに値が届くのを待って永久にブロックします。ループを終了できるよう、送信者は適切に破棄するかスコープを管理してください。
use std::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::channel();
let handle = thread::spawn(move || {
for i in 0..3 {
tx.send(i * 10).unwrap();
}
// tx dropped here, closing the channel
});
handle.join().unwrap();
let total: i32 = rx.iter().sum();
println!("Total: {}", total);
}ブロックしない読み取りのための try_recv
recv はブロックしますが、try_recv は Result をすぐに返します。メッセージが準備できていれば Ok(value) を、チャネルが空または切断されていれば Err を返します。ほかの処理も続ける必要があるイベントループで便利です。
use std::sync::mpsc;
use std::thread;
use std::time::Duration;
fn main() {
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
thread::sleep(Duration::from_millis(50));
tx.send("ready").unwrap();
});
loop {
match rx.try_recv() {
Ok(msg) => { println!("{}", msg); break; }
Err(_) => println!("waiting..."),
}
thread::sleep(Duration::from_millis(20));
}
}独自型を送信する
Send である任意の型をチャネル経由で送信できます。独自の構造体や列挙型も対象です。列挙型は、ワーカープロトコルで異なる種類のメッセージをモデル化するのに適しています。
use std::sync::mpsc;
use std::thread;
enum Job {
Print(String),
Add(i32, i32),
}
fn main() {
let (tx, rx) = mpsc::channel();
thread::spawn(move || {
tx.send(Job::Print(String::from("hi"))).unwrap();
tx.send(Job::Add(2, 3)).unwrap();
});
for job in rx {
match job {
Job::Print(s) => println!("print: {}", s),
Job::Add(a, b) => println!("add: {}", a + b),
}
}
}sync_channel とバックプレッシャー
mpsc::sync_channel(n) は、バッファサイズが n の有界チャネルを作成します。バッファが満杯になると、空きができるまで send はブロックします。これによりバックプレッシャーが働き、高速な生産者が低速な消費者を圧倒するのを防げます。
sync_channel(0)はランデブーチャネルです。send と recv が同時に成立する必要があります。
use std::sync::mpsc;
use std::thread;
fn main() {
let (tx, rx) = mpsc::sync_channel(2);
thread::spawn(move || {
for i in 1..=4 {
tx.send(i).unwrap();
println!("sent {}", i);
}
});
for v in rx {
println!("got {}", v);
}
}シンプルなワーカーパターン
チャネルはプロデューサー/コンシューマーパターンで特に力を発揮します。1つのスレッドが作業項目を生成し、別のスレッドがそれを受け取って処理します。ここでは、メインスレッドが数値を生成し、ワーカーがそれぞれを二乗して、2つ目のチャネル経由で結果を返します。
use std::sync::mpsc;
use std::thread;
fn main() {
let (job_tx, job_rx) = mpsc::channel();
let (res_tx, res_rx) = mpsc::channel();
thread::spawn(move || {
for n in job_rx {
res_tx.send(n * n).unwrap();
}
});
for n in 1..=4 {
job_tx.send(n).unwrap();
}
drop(job_tx);
for r in res_rx {
println!("square: {}", r);
}
}理解度チェック
mpscチャネルの理解度を確認しましょう。
まとめ
mpscチャネルの基本を学びました。
mpsc::channel()は(Sender, Receiver)のペアを返します。sendは値をチャネルに渡し、recvは値を取り出すまでブロックします。- Receiverを反復処理すると、すべてのsenderがdropされるまでメッセージを受信します。
- 複数のプロデューサーを使うには、
Senderをcloneします。 try_recvはノンブロッキングで、sync_channel(n)は容量制限付きのバックプレッシャーを追加します。
チャネルを使うと、メモリを共有するのではなく所有権を渡すことで、スレッド間で安全にデータを共有できます。
よくある質問
「mpscチャネル」レッスンは無料ですか?
はい。「mpscチャネル」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、Learn Rust Codingコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 Learn Rust Codingコースには全4レッスンが含まれています。
「mpscチャネル」で何を学びますか?
スレッド間で送信します ブラウザで直接実行するハンズオンコードでLearn Rust Codingを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
Learn Rust Codingを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのLearn Rust Codingは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン1/4です。
「mpscチャネル」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このLearn Rust Codingレッスンでコードを書いて実行できますか?
はい。すべてのLearn Rust Codingレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。