導入
「データを作る側(生産者)」と「処理する側(消費者)」を、安全に非同期でつなぎたい——そんなときの定番が System.Threading.Channels です。スレッド安全なキューのように振る舞い、生産と消費の速度差を吸収します。ブラウザ内ランナーには含まれないため、ここではコード例で読みます。
説明
flowchart LR
P["生産者<br/>Writer.WriteAsync"] --> C["Channel<br/>(スレッド安全なキュー)"]
C --> R["消費者<br/>await foreach Reader"]
style C fill:#e1f5fe
Channel.CreateUnbounded<T>() で作り、Writer に書き、Reader から読みます。消費側は await foreach で「届いた分から」処理できます。
サンプル
using System.Threading.Channels;
var channel = Channel.CreateUnbounded<int>();
// 生産者: 別タスクで数値を流し込み、終わったら Complete で閉じる
var producer = Task.Run(async () =>
{
for (int i = 1; i <= 3; i++)
{
await channel.Writer.WriteAsync(i);
}
channel.Writer.Complete(); // もう書かない、の合図
});
// 消費者: 届いた分から読み出す
await foreach (var item in channel.Reader.ReadAllAsync())
{
Console.WriteLine($"受信: {item}");
}
await producer;
Channel<T>はスレッド安全なキュー。生産者と消費者を非同期に橋渡しするWriter.WriteAsyncで書き、Reader.ReadAllAsyncをawait foreachで読むWriter.Complete()で「これ以上来ない」を伝える(消費側のループが終わる)- 容量制限つきの
CreateBoundedなら、速すぎる生産者に待ちをかけられる(背圧)
演習
手元の.NETプロジェクトで上のコードを動かし、生産者の WriteAsync の前に await Task.Delay(100) を入れて、消費側が「届いた分から」処理する様子を観察してみましょう。
まとめ
Channel<T>は生産者・消費者を安全につなぐ非同期キューawait foreach+ReadAllAsyncで届いた分から処理するCompleteで終了を伝え、Boundedで背圧をかけられる
次回: 複数のCPUで一気に処理する並列処理です(読み物)。