CodeGym /课程 /C# SELF /队列和栈 Producer-Consumer

队列和栈 Producer-Consumer

C# SELF
第 58 级 , 课程 1
可用

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); // 给消费者时间取完剩余项

注意,在这个简单示例里用标志 producersDoneThread.Sleep 来模拟完成。在真实应用中,更可靠的完成同步通常使用 CancellationTokenSourceBlockingCollection<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、表达式语法解析等)。

比较和选择合适的集合

集合 顺序 优点 典型场景
ConcurrentQueue
FIFO (先来先走) 按到达顺序提供公平的处理 任务队列、日志记录、处理入站请求、事件总线
ConcurrentStack
LIFO (后进先出) 快速访问最近添加的元素 操作历史(Undo/Redo)、对象池、遍历算法(DFS)

在 producer-consumer 场景里,选择 ConcurrentQueue 还是 ConcurrentStack 完全取决于你需要的处理顺序。两者都提供开箱即用的高性能和线程安全,让你不用操心手动同步,可以专注构建可扩展的多线程系统。

评论
TO VIEW ALL COMMENTS OR TO MAKE A COMMENT,
GO TO FULL VERSION