1. 介绍
先从基础开始。队列 是一个基本的数据结构,遵循 FIFO (First-In, First-Out) 原则,也就是“先来先走”。想象超市的排队:谁先排队,谁先被服务。
在多线程编程里,生产者-消费者 (Producer-Consumer) 模式是最常见也最有用的模式之一。
- 生产者 (Producers) 是生成数据或任务并把它们放到公共队列里的线程或应用部分。他们“生产”工作。
- 消费者 (Consumers) 是从队列取数据或任务并处理它们的线程或应用部分。他们“消费”工作。
这个模式能帮你控制数据流,把系统组件解耦(生产者不需要知道谁会如何处理数据),让应用更有响应性,并平滑地把负载分配到多个线程上。
示例:ConcurrentQueue — 添加和取出
看看如何向 ConcurrentQueue<T> 添加和取出元素。
using System.Collections.Concurrent;
ConcurrentQueue<string> tasks = new ConcurrentQueue<string>();
// 添加元素(生产者)
tasks.Enqueue("下载文件");
tasks.Enqueue("处理图像");
Console.WriteLine($"队列中的任务数: {tasks.Count}"); // 输出: 队列中的任务数: 2
// 取出元素(消费者)
if (tasks.TryDequeue(out string task1))
{
Console.WriteLine($"完成任务: {task1}"); // 输出: 完成任务: 下载文件
}
if (tasks.TryDequeue(out string task2))
{
Console.WriteLine($"完成任务: {task2}"); // 输出: 完成任务: 处理图像
}
if (!tasks.TryDequeue(out string emptyTask))
{
Console.WriteLine("队列为空,没有新任务。"); // 输出: 队列为空,没有新任务。
}
基础操作:Enqueue(), TryDequeue()
Enqueue(T item): 用来把元素添加到队列尾部。这个操作是线程安全的。你可以在 10 个不同线程同时调用 Enqueue,所有元素都会被正确添加。
TryDequeue(out T item): 用来尝试从队列头部取出元素。对消费者来说这是关键方法。如果成功取出元素会返回 true(值会传到输出参数 item),如果队列为空则返回 false。重要的是,TryDequeue 在队列为空时不会阻塞线程。
2. TryDequeue() 的重要性和操作的原子性
TryDequeue() 不仅仅是方便;它对正确的线程安全工作是关键的。它是原子性的:检查队列是否为空和实际取出元素是一整个不可分割的操作。
如果我们有单独的方法 IsEmpty(检查队列是否为空)和 Dequeue(取出元素),那么在两次调用之间其它线程可能会把队列清空。结果你的 Dequeue 会抛异常或返回错误数据。TryDequeue 完全避免了这种情况。
示例:多线程的 Producer-Consumer
这里我们启动两个生产者线程和一个消费者线程。
using System.Collections.Concurrent;
using System.Threading;
using System.Threading.Tasks;
ConcurrentQueue<int> dataQueue = new ConcurrentQueue<int>();
bool producersDone = false; // 用于向消费者发信号的标志
void Producer(int start, int count)
{
for (int i = 0; i < count; i++)
{
dataQueue.Enqueue(start + i);
Console.WriteLine($"[P] 添加: {start + i}");
Thread.Sleep(10);
}
}
void Consumer()
{
while (!producersDone || dataQueue.Count > 0) // 只要有数据或生产者还在工作就继续
{
if (dataQueue.TryDequeue(out int item))
{
Console.WriteLine($"[C] 处理: {item}");
}
else
{
Thread.Sleep(50); // 队列空时等待
}
}
Console.WriteLine("[C] 已完成工作。");
}
// 在 Main 中运行示例:
// Task.Run(() => Producer(1, 5));
// Task.Run(() => Producer(100, 5)); // 第二个生产者
// Task.Run(() => Consumer());
// Thread.Sleep(600); // 给线程一些时间工作
// producersDone = true; // 通知生产者已完成
// Thread.Sleep(200); // 给消费者时间取完剩余项
注意,在这个简单示例里用标志 producersDone 和 Thread.Sleep 来模拟完成。在真实应用中,更可靠的完成同步通常使用 CancellationTokenSource 或 BlockingCollection<T>。
ConcurrentQueue<T> 非常适合下面的场景:
- 处理顺序很重要(FIFO)。
- 多个线程同时添加元素,和/或多个线程同时取出元素。
- 需要高性能且不想手动管理锁。
3. 用于 producer-consumer 的栈 (LIFO)
栈 是另一种基础数据结构,遵循 LIFO (Last-In, First-Out) 原则,也就是“后进先出”。想象一摞盘子:你总是拿最上面的,新的盘子也放在最上面。
ConcurrentStack<T> 和 ConcurrentQueue<T> 一样是线程安全的,也可以用于“生产者-消费者”模式,但处理顺序是相反的。
示例:ConcurrentStack — 添加和取出
using System.Collections.Concurrent;
ConcurrentStack<string> commandStack = new ConcurrentStack<string>();
// 添加命令(生产者)
commandStack.Push("选中文本");
commandStack.Push("更改字体");
commandStack.Push("保存文档");
Console.WriteLine($"栈中的命令数: {commandStack.Count}"); // 输出: 栈中的命令数: 3
// 取出命令(消费者)
if (commandStack.TryPop(out string cmd1))
{
Console.WriteLine($"撤销命令: {cmd1}"); // 输出: 撤销命令: 保存文档
}
if (commandStack.TryPop(out string cmd2))
{
Console.WriteLine($"撤销命令: {cmd2}"); // 输出: 撤销命令: 更改字体
}
if (!commandStack.TryPop(out string emptyCmd))
{
Console.WriteLine("命令栈为空。"); // 输出: 命令栈为空。
}
4. 基础操作:Push(), TryPop()
Push(T item): 用来把元素添加到栈顶。操作是线程安全的。
TryPop(out T item): 用来尝试从栈顶取出元素。如果成功返回 true,如果栈为空返回 false。和 TryDequeue 一样,这是原子操作,可以防止竞态条件。
示例:用 ConcurrentStack 做对象池
栈很适合实现对象池:拿出-使用-放回。
using System.Collections.Concurrent;
class Connection { /* 简单占位 */ public Guid Id { get; } = Guid.NewGuid(); }
ConcurrentStack<Connection> connectionPool = new ConcurrentStack<Connection>();
// 用初始连接填充池
for (int i = 0; i < 3; i++)
{
connectionPool.Push(new Connection());
}
Console.WriteLine($"池里的连接数: {connectionPool.Count}"); // 输出: 池里的连接数: 3
void UseConnection()
{
if (connectionPool.TryPop(out Connection conn))
{
Console.WriteLine($"[池] 使用连接: {conn.Id}");
// 模拟使用连接
Thread.Sleep(50);
connectionPool.Push(conn); // 放回池里
Console.WriteLine($"[池] 归还连接: {conn.Id}. 池中: {connectionPool.Count}");
}
else
{
Console.WriteLine("[池] 连接池为空。创建新连接。");
// 通常在这里创建新连接,如果池为空
connectionPool.Push(new Connection());
}
}
// 在 Main 中运行示例:
Task.Run(() => UseConnection());
Task.Run(() => UseConnection());
Task.Run(() => UseConnection());
Thread.Sleep(500);
在这个例子里,多个线程可以安全地从共享池取连接并归还。
5. 应用示例和与 ConcurrentQueue 的比较
ConcurrentStack<T> 适用于:
- 当 LIFO 顺序很重要(例如“撤销”操作的历史)。
- 需要快速访问最近添加的元素(通常在 CPU 缓存中是“热”的)。
- 实现基于栈的算法(深度优先遍历 DFS、表达式语法解析等)。
比较和选择合适的集合
| 集合 | 顺序 | 优点 | 典型场景 |
|---|---|---|---|
|
FIFO (先来先走) | 按到达顺序提供公平的处理 | 任务队列、日志记录、处理入站请求、事件总线 |
|
LIFO (后进先出) | 快速访问最近添加的元素 | 操作历史(Undo/Redo)、对象池、遍历算法(DFS) |
在 producer-consumer 场景里,选择 ConcurrentQueue 还是 ConcurrentStack 完全取决于你需要的处理顺序。两者都提供开箱即用的高性能和线程安全,让你不用操心手动同步,可以专注构建可扩展的多线程系统。
GO TO FULL VERSION