0Pricing
C# Academy · レッスン

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()を呼び出すと、何が起こりますか?

まとめ

まとめ:上限付きの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フィードバックを取得できます。ローカル設定は不要です。

このコースのすべてのレッスン

  1. Parallel.ForEach、PLINQ
  2. Channelsによるプロデューサー/コンシューマー(概要)
  3. スループットとレイテンシーのトレードオフ
← C# Academyに戻る