使用 Channels 实现生产者/消费者(概览)
使用 BlockingCollection 建模类似 Channels 的生产者/消费者队列:有界缓冲区(背压)、多个消费者和优雅完成。
使用 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 实现平稳关闭。
用 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 反馈 — 无需本地设置。
此课程中的所有课时
- Parallel.ForEach 与 PLINQ
- 使用 Channels 实现生产者/消费者(概览)
- 吞吐量与延迟的权衡