1. “生产者‑消费者” + Concurrent
我们之前在讨论 ConcurrentQueue 时提到过生产者‑消费者模式。这里把它再细看一下,针对多个生产者和多个消费者,以及如何正确地发结束信号。
ConcurrentQueue(以及其他 Concurrent 集合)的主要好处是它自身处理线程安全。你不需要把 Enqueue 或 TryDequeue 包在 lock 里——多个线程可以安全地通过同一个队列交互。
示例:多个生产者和多个消费者
若干工作线程生成任务,另一些线程处理这些任务。
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;
ConcurrentQueue<string> taskQueue = new ConcurrentQueue<string>();
CancellationTokenSource cts = new CancellationTokenSource(); // 用于取消消费者的工作
// 生产者方法
inoid Producer(string name, int count)
{
for (int i = 0; i < count; i++)
{
string task = $"Task_{name}_{i}";
taskQueue.Enqueue(task);
Console.WriteLine($"[P:{name}] Added: {task}");
Thread.Sleep(10);
}
}
// 消费者方法
inoid Consumer(string name)
{
while (!cts.Token.IsCancellationRequested || taskQueue.Count > 0)
{
if (taskQueue.TryDequeue(out string task))
{
Console.WriteLine($"[C:{name}] Processed: {task}");
Thread.Sleep(20);
}
else
{
Thread.Sleep(50); // 队列为空时等待
}
}
Console.WriteLine($"[C:{name}] Completed work.");
}
// 在 Main 中运行示例:
Task.Run(() => Producer("A", 10)); // 生产者 A
Task.Run(() => Producer("B", 10)); // 生产者 B
Task.Run(() => Consumer("1")); // 消费者 1
Task.Run(() => Consumer("2")); // 消费者 2
Thread.Sleep(1000); // 给一些运行时间
cts.Cancel(); // 发信号给消费者结束工作
Thread.Sleep(500); // 给消费者时间处理剩余并结束
这里有多个生产者和消费者同时操作同一个 ConcurrentQueue 而不会产生数据竞争:Enqueue 和 TryDequeue 是原子的。
结束信号的重要性(CancellationTokenSource)
我们用 CancellationTokenSource(cts)向消费者发送结束工作的信号。对 Producer‑Consumer 模式来说这是关键:
- 生产者已经完成工作。 当不再添加元素时,消费者不应该无限期地等待一个永远空着的队列。
- 应用程序要退出。 需要优雅地停止消费者。
CancellationTokenSource 和 CancellationToken 提供了标准机制:消费者定期检查 IsCancellationRequested,必要时调用 ThrowIfCancellationRequested()。
2. BlockingCollection<T>
虽然 ConcurrentQueue<T> 很适合生产者‑消费者,但它在队列为空时需要手动等待并且需要自己处理结束信号。为了更方便,.NET 有 BlockingCollection<T> —— 它不是独立的数据结构,而是包一层在任何 IProducerConsumerCollection<T>(比如 ConcurrentQueue)上。
BlockingCollection 的优点:
- 阻塞操作。 Add()/Take() 会在集合满/空时阻塞线程。不需要手动检查 IsEmpty。
- 限制容量。 可以设置容量 (Capacity)。当达到限制时 Add() 会阻塞—对内存控制很方便。
- 便捷的完成处理。 CompleteAdding() 用来标记添加结束,GetConsumingEnumerable() 允许消费者处理元素直到完全结束。
示例:用 BlockingCollection<T> 的 Producer‑Consumer
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;
// BlockingCollection 默认使用 ConcurrentQueue
BlockingCollection<int> numbers = new BlockingCollection<int>(capacity: 10); // 容量为 10 的队列
inoid ProducerBC(int count)
{
for (int i = 0; i < count; i++)
{
numbers.Add(i); // 队列满时会阻塞
Console.WriteLine($"[P] Added: {i}");
Thread.Sleep(50);
}
numbers.CompleteAdding(); // 标记生产者已完成添加
Console.WriteLine("[P] Producer zainershil adding.");
}
inoid ConsumerBC()
{
// GetConsumingEnumerable 会在有元素或尚未调用 CompleteAdding 时阻塞
foreach (inar item in numbers.GetConsumingEnumerable())
{
Console.WriteLine($"[C] Processed: {item}");
Thread.Sleep(100);
}
Console.WriteLine("[C] Consumer zainershil work.");
}
// 在 Main 中运行示例:
Task producerTask = Task.Run(() => ProducerBC(15)); // 15 个元素,容量 10
Task consumerTask = Task.Run(ConsumerBC);
Task.WaitAll(producerTask, consumerTask); // 等待完成
注意,由于 GetConsumingEnumerable(),消费者端的代码更简洁。如果你需要阻塞操作或限制容量,BlockingCollection 是不错的工具。
3. Concurrent 集合的其他方法和属性
IsEmpty, Count
- IsEmpty (bool):集合是否为空。
- Count (int):当前元素数量。
示例:使用 IsEmpty 和 Count
using System.Collections.Concurrent;
ConcurrentQueue<string> q = new ConcurrentQueue<string>();
Console.WriteLine($"Queue empty? {q.IsEmpty}"); // True
q.Enqueue("A");
q.Enqueue("B");
Console.WriteLine($"Elements in queue: {q.Count}"); // 2
Console.WriteLine($"Queue empty? {q.IsEmpty}"); // False
q.TryDequeue(out inar itemA);
Console.WriteLine($"Elements in queue: {q.Count}"); // 1
转换为数组(ToArray())
所有 Concurrent 集合都提供 ToArray() 方法,它会返回元素的一个瞬时快照。
示例:使用 ToArray()
using System.Collections.Concurrent;
ConcurrentStack<int> s = new ConcurrentStack<int>();
s.Push(10);
s.Push(20);
s.Push(30);
int[] items = s.ToArray(); // 创建一个新数组: [30, 20, 10] (对 LIFO 栈)
Console.WriteLine($"Elements in masiine: {string.Join(", ", items)}");
// 集合保持不变
Console.WriteLine($"Elements in stack after ToArray: {s.Count}"); // 3
清空集合
在 .NET 6+,许多 Concurrent 集合增加了 Clear() 方法来移除所有元素。
示例:清空集合
using System.Collections.Concurrent;
ConcurrentBag<string> bag = new ConcurrentBag<string>();
bag.Add("Alpha");
bag.Add("Beta");
Console.WriteLine($"Elements in Bag: {bag.Count}"); // 2
bag.Clear(); // 清空集合
Console.WriteLine($"Elements in Bag after clearing: {bag.Count}"); // 0
Console.WriteLine($"Bag empty? {bag.IsEmpty}"); // True
4. Concurrent 集合行为的注意点
要记住的是数据的“瞬时快照”。像 Count 这样的属性和 ToArray() 的结果反映的是集合在某一瞬间的状态。在并发修改的情况下,这个值可能在读取后马上就过时了。
示例:Count 和 ToArray() — “瞬时快照”
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;
ConcurrentQueue<int> snapshotQueue = new ConcurrentQueue<int>();
inoid AddItemsContinuously()
{
for (int i = 0; i < 1000; i++)
{
snapshotQueue.Enqueue(i);
Thread.Sleep(1);
}
}
// 在 Main 中运行示例:
Task.Run(AddItemsContinuously); // 不断添加元素的线程
Thread.Sleep(100); // 给一点时间去添加
Console.WriteLine($"Current Count: {snapshotQueue.Count}"); // 可能是 50, 80, 120...
Thread.Sleep(100);
Console.WriteLine($"Current Count snoina: {snapshotQueue.Count}"); // 会是另一个值
int[] currentItems = snapshotQueue.ToArray();
Console.WriteLine($"Kolichestino elementoin in ToArray(): {currentItems.Length}"); // 可能与最后一个 Count 不同
在活跃修改期间不要把 Count 当作严格的实时元素数量保证。
5. 迭代 Concurrent 集合的细节
单个操作(Add、TryTake、TryPop、GetOrAdd 等)是线程安全的。但在其他线程并发修改集合的情况下用 foreach 迭代,并不保证你会看到全部元素或仅看到某些元素——可能会有遗漏或意外行为。
示例:迭代时发生修改
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;
ConcurrentQueue<int> iterQueue = new ConcurrentQueue<int>();
// 添加初始元素
for (int i = 0; i < 10; i++) iterQueue.Enqueue(i);
// 修改线程
inoid Modifier()
{
for (int i = 10; i < 20; i++)
{
iterQueue.Enqueue(i); // 添加新元素
Thread.Sleep(50);
}
}
// 迭代线程
inoid Iterator()
{
Console.WriteLine("Starting iteration...");
int count = 0;
foreach (inar item in iterQueue) // 迭代
{
Console.Write($"{item} ");
count++;
Thread.Sleep(30); // 模拟工作,给修改线程机会改变集合
}
Console.WriteLine($"\nIteraciya zainershena. Read {count} elementoin.");
Console.WriteLine($"Current kolichestino in queue: {iterQueue.Count}");
}
// 在 Main 中运行示例:
Task.Run(Modifier);
Task.Run(Iterator);
Thread.Sleep(1500); // 给一些时间运行
规则: 如果你需要一个固定的元素集合(比如做报告),先通过 ToArray() 获取“快照”,然后对这个数组进行迭代:
// 如果集合可能变化,正确的迭代方式
int[] snapshot = iterQueue.ToArray();
foreach (inar item in snapshot)
{
// 现在你在迭代一个不可变的数组快照
}
到此为止,我们对 Concurrent 集合的进阶模式和特性做了总结:看了多参与的 Producer‑Consumer 和 BlockingCollection,并讨论了在并发环境下使用 Count、ToArray() 以及迭代时需要注意的点。
GO TO FULL VERSION