Channelsによるプロデューサー/コンシューマー(概要)
BlockingCollection を使ってChannelsのようなプロデューサー/コンシューマーキューをモデル化します。容量制限付きバッファー(バックプレッシャー)、複数のコンシューマー、正常な完了を扱います。
「Channelsによるプロデューサー/コンシューマー(概要)」はCoddyKit上の無料C# Academyレッスンです。 これはレッスン2/3です。 下記で完全なレッスンを無料で読むことができます。その後、ブラウザ内の組み込みコードエディタと24時間対応のAIチューターでハンズオン演習できます。 これはC# Academy学習パスの一部であり、ウェブとCoddyKitアプリ全体で進捗が同期されます。 C# Academyコースには全3レッスンが含まれています。
概念と目標
目標:安全なプロデューサー/コンシューマーパイプラインを構築します。
- チャネルの考え方:プロデューサーとコンシューマーの間にあるキュー
- 上限付きの容量によってバックプレッシャーを実現します
- 複数のコンシューマーをサポートします
- 正常なシャットダウンにはCompleteAddingを使用します
単一プロデューサー/コンシューマー
基本的な流れは、プロデューサーがAddし、コンシューマーがGetConsumingEnumerableを反復処理し、最後にCompleteAddingを呼び出して終了します。
using System;
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;
public class Program
{
public static void Main(string[] args)
{
// Unbounded by default; simple demo
BlockingCollection<int> queue = new BlockingCollection<int>();
// Producer
Task producer = Task.Run(() =>
{
for (int i = 1; i <= 5; i++)
{
queue.Add(i); // enqueue
Console.WriteLine("Produced " + i);
Thread.Sleep(50); // simulate work
}
queue.CompleteAdding(); // signal no more items
});
// Consumer
Task consumer = Task.Run(() =>
{
foreach (int x in queue.GetConsumingEnumerable())
{
Console.WriteLine("Consumed " + x);
Thread.Sleep(80); // simulate processing
}
Console.WriteLine("Consumer done");
});
Task.WaitAll(producer, consumer);
}
}
バックプレッシャーの例
上限付きキューはバックプレッシャーを適用します。キューが満杯になると、コンシューマーが空きを作るまでAddがブロックされます。
using System;
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;
public class Program
{
public static void Main(string[] args)
{
// Bounded capacity introduces backpressure
BlockingCollection<int> queue = new BlockingCollection<int>(boundedCapacity: 2);
Task producer = Task.Run(() =>
{
for (int i = 1; i <= 5; i++)
{
queue.Add(i); // blocks when the buffer is full
Console.WriteLine("Produced " + i + " (count=" + queue.Count + ")");
}
queue.CompleteAdding();
});
Task consumer = Task.Run(() =>
{
foreach (int x in queue.GetConsumingEnumerable())
{
Console.WriteLine("Consumed " + x);
Thread.Sleep(120); // slower consumer -> producer will block sometimes
}
});
Task.WaitAll(producer, consumer);
}
}
ファンアウトコンシューマー
複数のコンシューマーがアイテムを競合して取得します(ファンアウト)。作業はスレッド間で自動的に分散されます。
using System;
using System.Collections.Concurrent;
using System.Linq;
using System.Threading;
using System.Threading.Tasks;
public class Program
{
public static void Main(string[] args)
{
BlockingCollection<int> queue = new BlockingCollection<int>(3);
Task prod = Task.Run(() =>
{
for (int i = 1; i <= 8; i++)
{
queue.Add(i);
Console.WriteLine("Produced " + i);
Thread.Sleep(30);
}
queue.CompleteAdding();
});
// Start 2 consumers that compete for items
Task[] consumers = Enumerable.Range(1, 2).Select(id => Task.Run(() =>
{
foreach (int x in queue.GetConsumingEnumerable())
{
Console.WriteLine("C" + id + " got " + x);
Thread.Sleep(100);
}
Console.WriteLine("C" + id + " done");
})).ToArray();
Task.WaitAll(consumers.Concat(new[] { prod }).ToArray());
}
}
正常な完了
キューが空になったときに正常に終了するには、CompleteAddingとGetConsumingEnumerableを使用します。
using System;
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;
public class Program
{
static void Producer(BlockingCollection<string> q)
{
string[] lines = { "A", "B", "C", "D" };
foreach (var s in lines)
{
q.Add(s);
Console.WriteLine("Produced " + s);
}
q.CompleteAdding(); // signal completion
}
static void Consumer(BlockingCollection<string> q)
{
try
{
foreach (var item in q.GetConsumingEnumerable())
{
Console.WriteLine("Consumed " + item);
Thread.Sleep(60);
}
Console.WriteLine("Consumer finished all items");
}
catch (InvalidOperationException)
{
// Thrown if taken after completion and empty; avoided by foreach above.
}
}
public static void Main(string[] args)
{
var queue = new BlockingCollection<string>(2);
Task p = Task.Run(() => Producer(queue));
Task c = Task.Run(() => Consumer(queue));
Task.WaitAll(p, c);
}
}
ヒントと注意点
ヒント:
- メモリの増加を防ぐため、上限付きバッファーを優先します。
- 完了時の競合を避けるため、GetConsumingEnumerableを使用します。
- 作業項目は小さく独立させ、共有状態を避けます。
- コンシューマー内でエラーをログに記録して処理し、処理が何も通知されずに停止するのを防ぎます。
CompleteAdding の目的
まとめ
まとめ:上限付きのBlockingCollectionを使用してチャネルを再現し、複数のコンシューマーにファンアウトし、正常なシャットダウンにはCompleteAddingを使用します。
よくある質問
「Channelsによるプロデューサー/コンシューマー(概要)」レッスンは無料ですか?
はい。「Channelsによるプロデューサー/コンシューマー(概要)」の完全なテキストはこのウェブで無料で読めます。インタラクティブに演習し(組み込みコードエディタと24時間対応のAIチューター)、C# Academyコースの残りをアンロックするには、CoddyKit PROにアップグレードしてください。 C# Academyコースには全3レッスンが含まれています。
「Channelsによるプロデューサー/コンシューマー(概要)」で何を学びますか?
BlockingCollection を使ってChannelsのようなプロデューサー/コンシューマーキューをモデル化します。容量制限付きバッファー(バックプレッシャー)、複数のコンシューマー、正常な完了を扱います。 ブラウザで直接実行するハンズオンコードでC# Academyを演習し、24時間対応のAIチューターがレッスンを進める中での質問に答えます。
C# Academyを始めるのに経験は必要ですか?
事前経験は必要ありません。CoddyKitのC# Academyは初級者から上級者向けに構成されているため、ここから始めるか最初から始めて、自分のペースで進むことができます。 これはレッスン2/3です。
「Channelsによるプロデューサー/コンシューマー(概要)」レッスンにはどのくらい時間がかかりますか?
ほとんどのCoddyKitレッスンは約5~10分かかります。各レッスンはコンパクトでインタラクティブなので、着実に進歩し、ウェブとアプリ全体で正確に前回の場所から再開できます。
このC# Academyレッスンでコードを書いて実行できますか?
はい。すべてのC# Academyレッスンに組み込みコードエディタが含まれているため、ブラウザでリアルコードを書いて実行し、即座のAIフィードバックを取得できます。ローカル設定は不要です。
このコースのすべてのレッスン
- Parallel.ForEach、PLINQ
- Channelsによるプロデューサー/コンシューマー(概要)
- スループットとレイテンシーのトレードオフ