资讯动态

C#并发编程:用ConcurrentQueue实现轻量级进程内消息队列

发布时间:2026/10/9 13:06:58 来源:尧图企业网站定制
前段时间我在重构一个工单处理服务遇到一件很典型的事用户提交一个工单后系统需要异步更新缓存、发送站内信、推送通知、写审计日志。最初这些操作全堆在请求线程里同步执行一个第三方通知接口超时整个请求就得跟着等最严重的时候一个接口平均响应时间从80毫秒飙到2秒多。我的第一反应是上RabbitMQ可评估下来发现只是单机部署、日请求量也就几十万条为这点量维护一套消息中间件实在不划算。后来我决定用C#内置的ConcurrentQueue自己封装一个进程内通信的消息队列把解耦和异步处理放在同一个进程里解决效果出乎意料地好。如果你也有类似场景——不想上重中间件、不想引入额外的网络依赖、只需要在单个进程内做生产者和消费者的解耦这篇文章应该能直接帮到你。我会从ConcurrentQueue的实现原理讲起手把手带大家实现一个支持背压、优雅停机、多消费者竞争消费的简易消息队列并分享我在实际项目中踩过的几个坑以及一组具体的压测数据。1. 进程内消息队列什么时候该自己动手什么时候该上中间件1.1 中间件不是银弹先判断场景是否小先给结论如果你的应用是单进程部署且队列场景集中在异步化处理、任务分发、事件通知这一类进程内消息队列是优先级最高的方案。原因很简单无网络开销消息不经过TCP、不经过序列化反序列化传递的是一个对象引用延迟在微秒级无部署成本不需要启动额外服务、不需要配置连接、不需要监控消费者连接数无跨进程的分布式问题不需要消息确认、不需要visibility timeout、不需要消息落盘整个方案简化一大截。RabbitMQ、Kafka这类中间件的优势在跨进程、跨机器、持久化、多订阅者场景下才真正体现。如果你的消费者需要跑在另一台机器上或者消息丢了不能接受或者需要堆积几十万条消息慢慢消费那还是要老老实实用中间件。进程内队列再怎么做也没办法处理进程崩溃之后的消息恢复问题。我自己的判断标准是这样的如果生产者和消费者的生命周期必须同时在同一个进程里且消息生命周期很短秒级、量不大百万条以内/天进程内队列完全够用。如果消费者需要断线重连、需要消费历史消息、需要水平扩展到多个进程实例那就别强行造轮子了。1.2 进程内排队要解决的三个核心问题所谓消息队列核心行为就是生产者往里面放消息消费者从里面取消息。在进程内实现它需要处理好三件事线程安全生产者和消费者往往运行在不同线程上队列本身必须支持并发读写这是最基本要求生产消费调度队列为空时消费者不能空转有消息时消费者要被及时唤醒这相当于一个轻量级的事件通知机制生命周期管理程序退出时队列里的存量消息怎么处理正在处理中的消息要不要等它处理完消费者的线程要不要优雅回收这三个问题前两个是ConcurrentQueue能直接解决的第三个需要我们自己设计。很多人以为ConcurrentQueue就是一个线程安全的List生产消费就是while循环里取一下但其实把第三个问题处理好才算真正达到生产可用级别。1.3 为什么选用ConcurrentQueue而不是BlockingCollection或Channel在.NET生态里进程内队列其实有不止一个选择我先把它们的区别摆在前面方案底层机制适用场景ConcurrentQueueCAS无锁队列需要自己控制阻塞策略、背压逻辑BlockingCollection默认基于ConcurrentQueue提供阻塞边界标准的生产者消费者中规中矩Channel独立优化的异步队列高吞吐、异步生产消费现代做法我选择从ConcurrentQueue讲起一是因为它够底层能把无锁并发的工作机制讲透理解它之后再看其他框架的源码也容易二是因为它的封装非常薄适合我们在它上面叠加自己的策略比如按优先级拒绝、按容量做背压、按消息Key做路由。BlockingCollection虽然自带阻塞拿取但在多消费者场景下对信号量的精细控制不如自己写来得直接。Channel 性能确实更好不过它的核心设计是围绕async/await的对于一批习惯了同步消费者模型的代码来说反而是ConcurrentQueue的手动控制方案更容易理解和掌控。2. 解读ConcurrentQueue无锁并发队列的底层秘密2.1 无锁并不意味着不需要同步先澄清一个常见误区很多人看到ConcurrentQueue是无锁lock-free实现就以为它完全不需要线程同步这是不对的。它确实没有使用lock关键字但底层依赖CPU的原子指令——CASCompare-And-Swap以及内存屏障Memory Barrier来保证多线程下的可见性和一致性。这样说可能还是有点抽象我换个生活化的比喻。传统的lock方式就像公司只有一个茶水间谁进去接水得拿钥匙用完了交出来其他人必须在外面排队等着等待的人什么也干不了。而CAS方式则是每人带一个自己的水杯两个人都准备同时往一个开水桶里接水硬件层面保证如果桶里水位还是我刚看到的水位那我就接如果变了我就重新看再尝试。不需要让其他人停下手里的其他工作等你处理速度快很多。在C#中CAS操作对应的就是Interlocked.CompareExchange这一族方法ConcurrentQueue入队和出队的核心逻辑就是围绕这些原子操作展开。之所以费这么大劲不用锁是因为锁的代价远比很多人想象中高线程在等待锁时会被挂起涉及内核态上下文切换高并发场景下这部分开销可能比业务逻辑本身还大。而CAS在无竞争或低竞争时开销极低基本就是一条CPU指令的代价。2.2 分段链表的存储结构ConcurrentQueue在底层并不是一个简单的链表或者数组而是一个分段链表segment-linked list。我说个直观的类比它像一本由很多页纸装订成的记事本每一页segment默认容量通常是32个元素可以连续存放一批元素页与页之间用链连接。入队时在最后一页的空位上写入如果最后一页满了就追加一个新页出队时从第一页取等第一页取空了就整体摘除释放。这种设计兼顾了两个目标一方面数据在段内连续存放访问时缓存友好比纯链表中每个节点散布各处要快得多另一方面它不需要像动态数组那样频繁扩容搬移数据。每次操作的平均复杂度依然是O(1)但实际缓存效率和吞吐量要优于手动用QueueT加锁的方案。这一点在消息量特别大的场景下体会尤其明显。2.3 使用中容易踩的API误区代码层面ConcurrentQueue提供了Enqueue、TryDequeue、TryPeek、IsEmpty、Count、Clear这几个常用成员用法本身不难但有三个细节在实际项目里很容易翻车判断队列是否为空用IsEmpty而不是Count 0。虽然看起来等价但ConcurrentQueue的Count带有近似性质并发写入后再立即读取结果不一定精确而且在不必要的场景下读Count还要承担额外的内存屏障开销。唯一推荐使用Count的场景是你明确知道当前没有并发或者你只需要一个大概大小。TryDequeue和TryPeek都是非阻塞的。队列为空时它们不会等待直接返回false。这意味着如果你在循环里快速调用TryDequeue而不做其他等待你的CPU会被白白烧掉——这是高CPU占用故障最常见的原因之一。foreach枚举是弱一致性的。用foreach遍历ConcurrentQueue不代表你拿到了一个稳定快照遍历过程中其他线程的入队出队可能体现出来也可能体现不出来。官方文档明确说不能依赖枚举做精确快照。这几个点后面都会反复用到尤其是TryDequeue的非阻塞特性以及为什么必须搭配信号量才能构建一个真正的阻塞式队列。3. 手工搭建生产者-消费者骨架队列 推送通知3.1 第一版轮询方案不推荐用于生产先上一版最朴素的实现。逻辑上完全正确但实际在真实环境里跑起来会发现明显问题。var queue new ConcurrentQueueOrderMessage(); // 生产者线程 Task.Run(() { for (int i 0; i 1000; i) { queue.Enqueue(new OrderMessage { Id i, CreatedAt DateTime.UtcNow }); Thread.Sleep(10); } }); // 消费者线程 Task.Run(() { while (true) { if (queue.TryDequeue(out var msg)) { Process(msg); } } });这个版本的问题有两个第一消费者线程内不断TryDequeue空转即使队列里一条消息都没有CPU占用率也是满的第二没有退出机制程序结束只能强杀线程。这两点在真实项目里都不可接受。尤其是CPU空转我见过一位同事在生产环境用这种模式处理TCP接收缓冲上线后单核CPU直接跑到100%排查了半天才发现是轮询坑了自己。3.2 第二版ConcurrentQueue SemaphoreSlim 实现阻塞式队列更优雅的做法是给队列搭配一个信号量。SemaphoreSlim是一个轻量级线程同步原语它的Count表示当前可用的资源数量Wait会阻塞线程直到Count大于0Release则让Count加1并唤醒一个等待者。利用这个机制可以设计出非常顺滑的生产消费模型入队一条消息就Release一次消费者取消息前先Wait一次。这样消费者在没有消息时会真正进入休眠状态而不是空转一旦有消息入队系统会唤醒一个正在等待的消费者。最关键的一点是SemaphoreSlim.Wait接收一个CancellationToken消费者可以在等待期间响应取消做到了等得到消息也停得下来。public sealed class BlockingMessageQueueT { private readonly ConcurrentQueueT _inner new ConcurrentQueueT(); private readonly SemaphoreSlim _signal new SemaphoreSlim(0); public int Count _inner.Count; public void Enqueue(T item) { _inner.Enqueue(item); _signal.Release(); } public T Dequeue(CancellationToken token) { _signal.Wait(token); // 阻塞等待直到有消息或取消 if (_inner.TryDequeue(out var item)) { return item; } // 单消费者场景这里不会走到多消费者场景可能出现 // 具体原因在第六章的坑里细讲。 throw new InvalidOperationException(signal and queue inconsistent); } }这里第_signal.Wait(token)被取消时会抛OperationCanceledException正好被消费者宿主捕获作为退出信号。这样消费者线程的循环逻辑就变得非常精简等待消息、取到消息、处理消息三件事干干净净。3.3 配套的消费者宿主与取消机制有了核心队列还需要一个消费者宿主负责循环拉取、执行业务、响应取消。下面这段代码展示如何用一个后台任务承载消费者循环并在程序退出时做到优雅停止public sealed class ConsumerWorkerT { private readonly BlockingMessageQueueT _queue; private readonly CancellationTokenSource _cts new CancellationTokenSource(); private readonly ListTask _workers new ListTask(); public ConsumerWorker(BlockingMessageQueueT queue, int workerCount) { _queue queue; for (int i 0; i workerCount; i) { _workers.Add(Task.Run(() Loop(_cts.Token))); } } private void Loop(CancellationToken token) { try { while (true) { var item _queue.Dequeue(token); Process(item); // 实际业务处理 } } catch (OperationCanceledException) { // 收到取消信号退出循环 } } public void Stop() { _cts.Cancel(); Task.WaitAll(_workers.ToArray()); } }Stop方法先取消再等待所有Task结束保证消费者处理完当前消息后才整体退出。这就是一个非常精简但拿得出手的进程内消息队列骨架。下一步要考虑的是业务模型了多个消费者之间到底是竞争关系还是订阅关系。4. 进阶设计竞争消费与广播订阅两种模型4.1 竞争消费多个消费者瓜分消息前面实现的队列天然支持竞争消费——多个消费者线程同时调用Dequeue每条消息只会被其中一个消费者取走这正是任务分发场景需要的模式。比如你有8个消息处理Worker它们共享同一个队列各自取消息处理处理完再取下一条天然实现负载均衡。这种模式下要注意的是消息处理的幂等性每条消息只被一个消费者处理一次是在队列层保证的但如果Process方法内部抛了异常消息已经出队了这条消息就相当于丢了。这需要业务层自己做重试策略。我的习惯是在Process内部做有限次重试失败后把消息写入一个单独的死信队列本质就是另一个ConcurrentQueue同时触发告警。4.2 广播订阅每个订阅者一个队列另一类更常用的模型是广播订阅。同一类消息不只是一个人关心订单创建后积分系统要加积分、短信系统要发通知、统计系统要记日志。竞争消费的模式在这里就不行了因为一个消费者拿走消息后其他消费者就再也看不到这条消息。这时候的标准做法是每个订阅者维护一个自己的BlockingMessageQueue发布者把消息复制发给所有订阅者。核心代码大致是这样public sealed class PubSubHubT { private readonly ListBlockingMessageQueueT _subscribers new ListBlockingMessageQueueT(); public IDisposable Subscribe(BlockingMessageQueueT queue) { lock (_subscribers) { _subscribers.Add(queue); } return new Unsubscriber(_subscribers, queue); } public void Publish(T message) { ListBlockingMessageQueueT snapshot; lock (_subscribers) { snapshot _subscribers.ToList(); // 快照避免遍历时订阅者增删 } foreach (var subscriber in snapshot) { subscriber.Enqueue(message); } } private sealed class Unsubscriber : IDisposable { private readonly ListBlockingMessageQueueT _list; private readonly BlockingMessageQueueT _queue; public Unsubscriber(ListBlockingMessageQueueT list, BlockingMessageQueueT queue) { _list list; _queue queue; } public void Dispose() { lock (_list) { _list.Remove(_queue); } } } }两个容易踩的坑第一遍历订阅者集合时如果直接用foreach遍历订阅/退订操作会造成集合已修改异常所以必须做快照第二Publish返回时消息只是进了各个订阅者的队列不代表所有订阅者已经处理完毕——这是异步解耦的固有特性。如果某个订阅者消费速度远低于发布速度它的独立队列就会慢慢膨胀这就引出了第五章要讲的背压机制。4.3 消息体的线程安全设计队列本身线程安全不代表消息体也安全。如果你把同一个消息对象Publish到多个订阅者队列这个对象实际上是在多个消费者线程间共享的任何消费者修改了消息对象的属性其他消费者都会看到。这个坑极其隐蔽某个订阅者在处理消息时把msg.Status 1另一个订阅者消费时看到的已经是1而不是初始状态。我给自己定的规矩是广播订阅中的消息必须是不可变对象或者在Publish时给每个订阅者传入浅拷贝快照。要么消息类所有属性只读要么在Publish阶段就完成拷贝。竞争消费场景下相对安全但也要注意消息体内部的可变集合或引用类型不要被多个消费者交叉修改。5. 生产环境必不可少的三个细节背压、顺序、优雅停机5.1 背压让队列不要无限膨胀进程内队列最危险的问题不是并发错乱而是发布速度远大于消费速度时队列无限膨胀最终把内存撑爆。RabbitMQ有磁盘、有流控帮你缓冲进程内队列可没有这样的能力。所以生产环境里我总会给队列加一个容量上限超过上限时生产者要么被阻塞等待要么按策略丢弃。最简单的有界队列改造核心思路还是放在信号量上。引入两个信号量一个表示队列中已有多少消息可消费信号一个表示队列中还剩多少容量容量信号。入队前先capacitySem.Wait()入队后就signalSem.Release()出队前先signalSem.Wait()出队后capacitySem.Release()。一旦队列满到容量上限生产者就会阻塞在容量信号上从而自然实现背压。public sealed class BoundedBlockingQueueT { private readonly ConcurrentQueueT _inner new ConcurrentQueueT(); private readonly SemaphoreSlim _signal new SemaphoreSlim(0); private readonly SemaphoreSlim _freeCapacity; public BoundedBlockingQueue(int capacity) { _freeCapacity new SemaphoreSlim(capacity); } public bool TryEnqueue(T item, int timeoutMs) { if (!_freeCapacity.Wait(TimeSpan.FromMilliseconds(timeoutMs))) { return false; // 队列满按策略丢弃或降级 } _inner.Enqueue(item); _signal.Release(); return true; } public T Dequeue(CancellationToken token) { _signal.Wait(token); var ok _inner.TryDequeue(out var item); _freeCapacity.Release(); return item; } }这种背压方案的核心价值在于它让生产过快这个问题以一种可控的方式暴露出来。生产者在入队处等待超时时说明下游真的来不及处理了这时候正确的策略是让调用方感知到压力超时、降级日志、限流丢弃而不是让进程默默把内存吃光。这里有一个容易被忽略的细节如果_freeCapacity.Wait超时返回false就不能再执行入队逻辑上必须严格保持容量信号和实际容量一致否则后面会错乱。5.2 顺序保证多消费者会打乱你的消息顺序进程内消息队列本身是FIFO结构队列头部永远是最先入队的消息这没问题。但这只在你只有一个消费者、或者把入队顺序等同于全局处理顺序时才成立。一旦启用多个消费者线程两个线程同时取到相邻消息谁先处理完是不确定的全局处理顺序自然就被打乱了。如果需要保证严格顺序方案只有一个让需要保序的那一组消息永远只走同一个消费者。常见做法是按业务Key做哈希路由比如订单号、用户ID。把相同Key的消息归到同一个消费者关联的队列这样每个队列内部依然是FIFO整体看来同一业务Key的消息处理顺序就有保证了。如果你的场景是所有消息都必须全局严格有序那多消费者就是不可接受的设计老老实实用单消费者。绝大多数业务其实不需要全局严格有序能接受一定程度的乱序用它换吞吐量是完全划算的。5.3 优雅停机让存量消息处理完再关闭进程进程内队列的停机问题实际项目里被忽视得最厉害。很多人写的消费者循环是这样while (true) { var msg queue.Dequeue(); Process(msg); }这里的消费者线程是后台线程或者用了Task但没有在退出时调用WaitAll——进程一结束队列里积压的消息还在正在处理的消息直接被砍断。对于消息丢失不可接受的场景这就属于生产事故。我认为比较可靠的停机流程分三步先通知生产者停止入队外部请求不再往里发新消息然后取消消费者等待令牌消费者在Dequeue处收到OperationCanceledException退出如果消费者此刻正在Process某条消息会等Process完成后再进入下一次Dequeue下一次才检测到取消最后用Task.WaitAll(workers)等待所有消费者线程全部退出。此时队列里仍未消费的消息自然被丢弃业务方自行记录或告警但已经取出且正在处理的消息不会中断。配合前面BoundedBlockingQueue停机的关键代码就是Stop()里先Cancel再WaitAll。如果消费者数量多且单条任务执行时间长最好给WaitAll加一个超时时间或者结合Task.WhenAny做超时处理避免停机过程被一条慢任务无限卡死。6. 实际项目踩过的坑与性能实测6.1 空转消费导致CPU飙满这是我在真实项目里踩到的第一个坑。某个服务的探活模块用队列做缓冲消费者线程写成了无限TryDequeue空转上线后CPU直接跑满。排查过程并不复杂——用dotnet-counters看到上下文切换次数极高再抓dump看线程栈发现消费者线程全在ConcurrentQueue的TryDequeue里打转。解决办法就是现在这版代码的方式用SemaphoreSlim把消费者从轮询改成监听。这里我刻意保留一个明确的原则生产环境的消费者循环不应该在没有消息时空转哪怕你为了业务需要保持一个低频轮询比如100ms一次也不要无脑疯转。空转不仅浪费CPU还会给排查问题制造大量噪点。6.2 多消费者场景下的信号量错乱一次幽灵唤醒分享一个更隐蔽的坑。最初实现多消费者时我用的模式是while (true) { _signal.Wait(); if (_queue.TryDequeue(out var msg)) { Process(msg); } }表面看没什么问题但多消费者场景下信号量的Release和队列的Enqueue并不是一个原子操作生产者入队后Release此时有两个消费者在Wait上A被唤醒并成功拿到消息B也可能被同时唤醒但发现队列已空于是B在TryDequeue返回false后进入下一轮循环再次调用Wait阻塞。问题就出在这里——信号量的Count已经被B的Wait扣掉了但B并没有真正拿到消息这破坏了一个信号量计数对应一条真实消息的不变量。后续再有消息入队就只有一个消费者能被唤醒另外一个永远阻塞造成吞吐量减半甚至服务假死。所以我才用最开始的BlockingMessageQueue.Dequeue设计——先Wait再TryDequeue如果TryDequeue失败了就说明信号量和队列状态不一致必须抛异常暴露出来而不是静默吞掉。严谨的多消费者队列实现要么用BlockingCollection这种封装好的边界处理要么给每个消费者分配独立的信号量避免共享信号量的计数被稀释。如果你想坚持自己写一定要在一次Wait成功但TryDequeue失败之后重新Wait并且做好日志记录因为这大概率意味着入队和信号量的配对出现了设计漏洞。6.3 性能实测ConcurrentQueue对比加锁Queue最后给一组我在一台8核机器上做的简单压测数据仅供参考。测试内容是两百万条消息生产者和消费者各5个线程统计从入队到全部出队完成的总耗时方案总耗时约说明lock (lockObj) { _queue.Enqueue/Dequeue }3.2秒全程有锁竞争时线程阻塞开销大裸ConcurrentQueueT1.8秒无锁低竞争时吞吐优势明显自研SemaphoreSlim ConcurrentQueue1.9秒比裸ConcurrentQueue多了一点信号量开销但换来了阻塞和停机能力差距没有想象中巨大但能明确感受到ConcurrentQueue在低竞争时确实更快。而信号量方案的目的从来不是极限压测跑分而是要解决空转、停机和背压这些工程问题。并发场景下程序正确性和资源占用往往比那零点几秒的吞吐更重要。6.4 什么时候可以直接换成Channel或BlockingCollection如果你看完这篇文章觉得手动管理信号量太麻烦完全可以考虑直接用ChannelT。Channel在.NET的异步生态里是官方推荐的高吞吐方案内部做了大量优化支持异步读写、支持背压通过BoundedChannelOptions、支持多消费者代码量还比我们的方案少得多。我们这套ConcurrentQueue方案更适合的场景是你需要完全掌控入队出队策略比如自定义优先级、自定义丢弃规则或者项目里还有大量同步代码要接入。一个很实际的建议是先按这篇文章的思路自己实现一遍把并发模型和背压机制彻底摸透然后评估性能或代码量如果不够满意再切Channel。反过来如果你只是想要一个能跑的进程内队列我的观点也很明确——你不用自己造轮子但你应该看得懂轮子是怎么转的。这恰恰是高级这两个字的含金量所在。我在项目里从这套ConcurrentQueue方案起步后来一步步又加了死信队列、延迟队列、动态消费者伸缩但解决核心问题的思路始终没变过线程安全地移交消息异步地处理副作用。希望这篇实操记录能帮到正在琢磨类似方案的你。

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价 →
↑