C# Academy · 课时

使用 Channels 实现生产者/消费者(概览)

使用 BlockingCollection 建模类似 Channels 的生产者/消费者队列:有界缓冲区(背压)、多个消费者和优雅完成。

第 2 / 3 课8 个步骤

使用 Channels 实现生产者/消费者(概览) 是 CoddyKit 上的免费 C# Academy 课时。 这是第 2 节课,共 3 节。 你可以在下方免费阅读本课时的完整内容 — 然后在浏览器中使用内置代码编辑器和全天候 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 实现平稳关闭。

免费开始

用 AI 导师学习 C# — 免费

在浏览器中编写并运行真实代码,获得全天候 AI 导师的即时帮助,并在网页或应用中继续学习。

课程
93
课程
346

常见问题解答

「使用 Channels 实现生产者/消费者(概览)」课时是免费的吗?

是的 — 「使用 Channels 实现生产者/消费者(概览)」的完整文本可在网页上免费阅读。要进行交互式练习(内置代码编辑器和全天候 AI 导师)并解锁 C# Academy 课程的其余内容,请升级到 CoddyKit PRO。 C# Academy 课程共包含 3 节课。

「使用 Channels 实现生产者/消费者(概览)」这节课中我会学到什么?

使用 BlockingCollection 建模类似 Channels 的生产者/消费者队列:有界缓冲区(背压)、多个消费者和优雅完成。 你通过在浏览器中直接运行的动手代码来练习 C# Academy,全天候 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