System.Threading.Channels
Channel、ChannelWriter、ChannelReaderを使ってプロデューサー・コンシューマーパイプラインを構築し、バックプレッシャーを制御します。
「System.Threading.Channels」はCoddyKit上の無料C# Academyレッスンです。 これはレッスン2/4です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはC# Academy学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 C# Academyコースには全4レッスンが含まれています。
チャネルとは
System.Threading.Channels(.NET Core 3で導入)は、高性能でスレッドセーフなプロデューサー・コンシューマーキューを提供します。BlockingCollectionとは異なり、チャネルは完全に非同期でアロケーション効率も高いため、プロセス内パイプラインに適しています。
チャネルの作成
チャネルはファクトリを使って作成します。無制限(上限なし)または有界(容量に上限があり、バックプレッシャーが発生する)を選択します。ファクトリは、WriterとReaderを備えたChannel<T>を返します。
using System.Threading.Channels;
// Unbounded: unlimited capacity, no backpressure
var unbounded = Channel.CreateUnbounded<string>();
// Bounded: max 100 items; writer waits when full
var bounded = Channel.CreateBounded<string>(100);
// Bounded with drop-oldest strategy:
var dropping = Channel.CreateBounded<string>(new BoundedChannelOptions(50)
{
FullMode = BoundedChannelFullMode.DropOldest
});チャネルへの書き込み
プロデューサーはWriteAsync(有界チャネルが満杯の場合は待機)またはTryWrite(満杯の場合はすぐにfalseを返す)を使って項目を書き込みます。Writer.Complete()で完了を通知します。
var channel = Channel.CreateUnbounded<int>();
// Producer task
var producer = Task.Run(async () =>
{
for (int i = 0; i < 100; i++)
{
await channel.Writer.WriteAsync(i);
await Task.Delay(10);
}
channel.Writer.Complete(); // signal no more items
});チャネルからの読み取り
コンシューマーは、最も簡単な方法であるReadAllAsync()か、より細かく制御できるReadAsync/TryReadを使って読み取ります。ReadAllAsyncは、ライターがComplete()を呼び出すと完了します。
// Consumer task
var consumer = Task.Run(async () =>
{
await foreach (var item in channel.Reader.ReadAllAsync())
{
Console.WriteLine($"Processed: {item}");
// Naturally paced — waits for next item
}
Console.WriteLine("Channel completed");
});
await Task.WhenAll(producer, consumer);複数のプロデューサー
チャネルはスレッドセーフです。複数のプロデューサーがロックなしで同時に書き込めます。Writer.Complete()を呼び出すのは、すべてのプロデューサーが終了した後だけにしてください。
var channel = Channel.CreateUnbounded<WorkItem>();
var producers = Enumerable.Range(0, 4).Select(id =>
Task.Run(async () =>
{
for (int i = 0; i < 25; i++)
await channel.Writer.WriteAsync(new WorkItem(id, i));
}));
// Wait for all producers before completing the writer
await Task.WhenAll(producers);
channel.Writer.Complete();複数のコンシューマー(ファンアウト)
同じチャネルリーダーに対して複数のコンシューマータスクを実行すると、処理を並列化できます。各項目はちょうど1つのコンシューマーに配信されます(ブロードキャストではなく分割配信です)。
var channel = Channel.CreateBounded<WorkItem>(100);
// 4 parallel consumers
var consumers = Enumerable.Range(0, 4).Select(id =>
Task.Run(async () =>
{
await foreach (var item in channel.Reader.ReadAllAsync())
{
await ProcessItemAsync(item);
Console.WriteLine($"Consumer {id} processed {item.Id}");
}
}));
await Task.WhenAll(consumers);パイプラインパターン
チャネルを連結して処理パイプラインを構築します。各ステージが1つのチャネルから読み取り、項目を処理して、次のチャネルへ書き込みます。各ステージは自然なバックプレッシャーを伴って同時に実行されます。
// Stage 1: raw data
var stage1 = Channel.CreateBounded<string>(50);
// Stage 2: parsed
var stage2 = Channel.CreateBounded<ParsedRecord>(50);
// Stage 3: enriched
var stage3 = Channel.CreateBounded<EnrichedRecord>(50);
var parse = ParseStageAsync(stage1.Reader, stage2.Writer);
var enrich = EnrichStageAsync(stage2.Reader, stage3.Writer);
var persist = PersistStageAsync(stage3.Reader);
await Task.WhenAll(parse, enrich, persist);有界チャネルによるバックプレッシャー
有界チャネルはバックプレッシャーを自動的に適用します。チャネルが満杯になると、コンシューマーが項目を取り出すまでWriteAsyncがプロデューサーを一時停止します。手動でスロットリングするコードは必要ありません。
// Bounded channel: max 10 items
var channel = Channel.CreateBounded<string>(10);
// Fast producer
var producer = Task.Run(async () =>
{
for (int i = 0; i < 1000; i++)
{
// WriteAsync waits when channel has 10 items
await channel.Writer.WriteAsync($"item-{i}");
// Producer is naturally slowed to consumer speed
}
channel.Writer.Complete();
});
// Slow consumer
var consumer = Task.Run(async () =>
{
await foreach (var item in channel.Reader.ReadAllAsync())
{
await Task.Delay(50); // simulate slow processing
Console.WriteLine(item);
}
});チャネルでのエラー処理
Writer.Complete(exception)に例外を渡すと、待機中のすべてのリーダーへエラーを伝播できます。リーダーは、次にチャネルを読み取る際に例外を受け取ります。
var channel = Channel.CreateUnbounded<int>();
var producer = Task.Run(async () =>
{
try
{
for (int i = 0; i < 100; i++)
{
if (i == 50) throw new Exception("Producer failed at 50");
await channel.Writer.WriteAsync(i);
}
channel.Writer.Complete();
}
catch (Exception ex)
{
channel.Writer.Complete(ex); // propagate to reader
}
});
try
{
await foreach (var item in channel.Reader.ReadAllAsync())
Console.WriteLine(item);
}
catch (Exception ex)
{
Console.Error.WriteLine($"Channel error: {ex.Message}");
}実践例:バックグラウンドジョブキュー
ホステッドサービスとチャネルを使ったバックグラウンドジョブキューでは、HTTPリクエストがジョブをキューに追加し、バックグラウンドワーカーがジョブを1つずつ処理します。
public class JobQueue
{
private readonly Channel<Func<CancellationToken, Task>> _queue
= Channel.CreateBounded<Func<CancellationToken, Task>>(100);
public ChannelWriter<Func<CancellationToken, Task>> Writer => _queue.Writer;
public ChannelReader<Func<CancellationToken, Task>> Reader => _queue.Reader;
}
public class JobProcessor : BackgroundService
{
private readonly JobQueue _queue;
public JobProcessor(JobQueue q) => _queue = q;
protected override async Task ExecuteAsync(CancellationToken ct)
{
await foreach (var job in _queue.Reader.ReadAllAsync(ct))
await job(ct);
}
}確認問題
満杯の有界チャネルに対してプロデューサーがWriteAsyncを呼び出すと、どうなりますか。
まとめ:System.Threading.Channels
重要なポイント:
- チャネルは、高性能で非同期処理に適したプロデューサー・コンシューマーキューを提供します
- 無制限:上限なし。有界:バックプレッシャーまたはドロップ戦略を伴う有限容量です
- 複数のプロデューサーとコンシューマーを、ロックなしで安全に使用できます
- ReadAllAsync() + await foreach は、最も簡潔な消費パターンです
- チャネルをパイプラインに連結すると、ステージ単位の並列処理を実現できます
- Writer.Complete()(またはComplete(exception))を呼び出して、ストリームの終了を通知します
よくある質問
「System.Threading.Channels」レッスンは無料ですか?
はい。「System.Threading.Channels」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、C# Academyコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 C# Academyコースには全4レッスンが含まれています。
「System.Threading.Channels」で何を学びますか?
Channel、ChannelWriter、ChannelReaderを使ってプロデューサー・コンシューマーパイプラインを構築し、バックプレッシャーを制御します。 ブラウザで直接実行するハンズオンコードでC# Academyを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
C# Academyを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのC# Academyは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン2/4です。
「System.Threading.Channels」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このC# Academyレッスンでコードを書いて実行できますか?
はい。すべてのC# Academyレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- IAsyncEnumerableとawait foreach
- System.Threading.Channels
- ValueTaskとアロケーションの回避
- ConfigureAwaitとSynchronizationContext