拓十年匠心定制 · 商业建站与技术教学双线并行 咨询热线:400-886-1026 service@lmnt.cn
ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

基于ConcurrentQueue实现进程内消息队列:线程解耦与异步处理实战

基于ConcurrentQueue实现进程内消息队列:线程解耦与异步处理实战

1. 为什么进程内也需要消息队列

1.1 什么样的场景真正需要它

很多人一听到"消息队列",第一反应就是 RabbitMQ、Kafka 这一类的分布式中间件。但实际开发里,大量场景根本不需要跨进程、跨机器的消息分发,痛点仅仅出在同一个进程内的多个线程之间需要解耦通信。

举个例子,我做过一个 C# 上位机项目,设备通过串口不停上报数据,采集线程拿到原始报文后要做解析,解析结果要分发给 UI 线程刷新界面、分发给业务逻辑线程做告警判断、还要丢给日志线程记录。如果直接在采集线程里同步调用各个模块的方法,采集线程会被拖死——UI 刷新慢、日志写盘慢,全都变成串联依赖。这时候就需要一个进程内的消息管道:采集线程只管往管道里扔消息,其他模块各起各的消费线程去取。

再比如后台任务系统,一个订单处理流程包含校验、库存锁定、支付回调、通知推送好几个环节。如果全都用同步方法调用,任何一个环节出了问题,整个调用链都要跟着等待。用进程内队列把每个环节拆开,上游只管投递消息,下游异步处理,系统的响应速度和容错能力立刻就不一样了。

这类场景的共同特征是:不要求消息跨机器传递,也不要求持久化,但要解决线程之间的解耦、异步削峰和背压问题。这时候上 RabbitMQ 属于杀鸡用牛刀——部署依赖、网络开销、运维成本全上来了,而且数据还要序列化反序列化走一遍网络,延迟远高于进程内直接传递。用一个线程安全的队列,比如ConcurrentQueue<T>,就能以极低的成本解决 80% 的问题。

1.2 轻量级方案选型的权衡

选ConcurrentQueue<T>而不是其他方案,需要把几个候选放在一起看:

方案线程安全阻塞/非阻塞适用场景注意事项
ConcurrentQueue<T>是非阻塞多生产者、多消费者,追求吞吐量队列为空时需自旋或配合信号机制
BlockingCollection<T>是阻塞生产者消费者模型,天然支持限流底层默认用ConcurrentQueue,有阻塞等待能力
Channel<T>是异步非阻塞追求高性能异步流水线API 更现代,支持只读/只写视图,但需要.NET Core 3.0+
List<T>+lock手动加锁非阻塞简单场景,消息量小锁竞争激烈时性能急剧下降

我在项目里选择ConcurrentQueue<T>作为底层存储,原因有三点。第一,它是无锁设计,内部使用 CAS 操作,高并发下吞吐量比List<T>加 lock 高一个数量级;第二,它天然支持多生产者多消费者,不限制生产者和消费者的数量;第三,它没有阻塞机制,不会因为队列满而挂起生产者线程,这对于需要控制线程生命周期的场景更灵活。

但ConcurrentQueue<T>也有一个明显的短板:没有阻塞等待能力。消费者在队列为空时调用TryDequeue会立即返回false,如果消费者用while(true)循环去轮询,CPU 空转问题会很严重。这个问题我在后面的核心实现里给出解决方案,这也是这篇文章真正有价值的地方——不是简单告诉你ConcurrentQueue怎么用,而是怎么把它做成一个真正能上生产环境的进程内消息队列。

2. ConcurrentQueue 的核心机制与正确用法

2.1 底层实现:为什么它能做到线程安全

ConcurrentQueue<T>是 .NET 框架里少有的几个无锁数据结构之一。理解它的底层实现,对写出高性能代码大有帮助。

它的内部是由多个**存储段(Segment)**组成的链表,每个 Segment 是一个小数组,用来存放实际数据。队列头部和尾部分别由head和tail指针维护,入队操作通过Interlocked系列的无锁原子操作更新尾部指针,出队操作则通过原子操作更新头部指针。

// 简化示意,理解核心思想即可 internal class ConcurrentQueueSegment<T> { internal volatile T[] _array; // 数据存储 internal volatile int _head; // 段内头部索引 internal volatile int _tail; // 段内尾部索引 }

关键在于,入队和出队操作不会互相干扰——入队只修改尾部段,出队只修改头部段。多个生产者同时入队时,使用Interlocked.Increment竞争获取唯一的槽位索引,拿到索引后直接写入,不需要锁住整个队列。这就是为什么ConcurrentQueue<T>在并发写入场景下表现极佳。

注意:ConcurrentQueue<T>是弱一致性的。枚举器在遍历期间,如果有其他线程同时入队或出队,枚举结果可能包含部分已移除的元素或遗漏部分新加入的元素。这不算 bug,在设计上就是如此——换取的是一些场景下更高的并发性能。如果你需要强一致性的快照,应该先ToArray()或者加锁。

理解这个底层机制后,你就明白为什么官方文档建议:大量入队、少量出队的场景选ConcurrentQueue,大量出队、少量入队的场景考虑ConcurrentBag。因为出队操作在队列头部方向竞争,而ConcurrentQueue的头部段在回收旧段时需要处理内存释放,有一个近似"拆段"的开销。不过对绝大多数进程内消息队列场景,这个差异完全可以忽略。

2.2 API 使用中的关键细节

ConcurrentQueue<T>的核心 API 不多,就几个:

  • Enqueue(T item):入队,添加到队列尾部,永远不阻塞,不会抛异常(除非item为 null 且 T 是引用类型,实际不允许 null 入队)。
  • TryDequeue(out T result):尝试出队,成功返回true,队列空返回false,这个方法也是异步安全的。
  • TryPeek(out T result):尝试查看队首元素但不移除。
  • Count:获取队列中元素数量。
  • IsEmpty:判断队列是否为空,比Count == 0效率略高,但也是弱一致的。
  • Clear():清空队列。

其中有几个细节我踩过坑,值得展开说。

第一个坑:TryDequeue的返回值处理。很多新手写消费者循环时是这样写的:

while (queue.TryDequeue(out var msg)) { Process(msg); }

看起来没问题,但这个循环在队列为空时会立刻退出,如果此刻生产者还没开始投递消息,消费者就已经"收工"了。正确的方式是外层再包一层循环,并配合信号机制等待。

第二个坑:Count属性的弱一致性。在高并发下,Count返回的是一个近似值。如果某个线程刚Enqueue完,立刻在另一个线程读Count,可能读到的是旧值。假如你拿Count做积压告警,阈值判断就要留出余量,否则告警会抖动。

第三个坑:null不能入队。ConcurrentQueue<T>内部会阻止 null 元素入队,如果你需要传递空对象,应该用包装类或者定义消息基类。

我来写一个安全的消费者循环模板,这是后面整个消息队列的核心基础:

while (!_cancellationToken.IsCancellationRequested) { // 先尝试非阻塞出队 if (_queue.TryDequeue(out var message)) { Process(message); continue; } // 队列为空,等待信号,而不是空转轮询 try { _signal.WaitOne(_cancellationToken); } catch (OperationCanceledException) { break; } }

_signal是一个AutoResetEvent或者SemaphoreSlim,生产者入队后调用信号释放,消费者收到信号后才重新尝试出队。这样既避免了 CPU 空转,也保证了消息的实时性。后面的完整实现就是在这个模板上扩展的。

3. 从零搭建简易进程内消息队列

3.1 消息模型设计与整体架构

动手写代码之前,先把架构想清楚。一个完整的进程内消息队列,最少包含四个部分:

  1. 消息本身(Message):承载业务数据和路由信息。
  2. 队列管理器(MessageQueue):内部持有ConcurrentQueue<T>,提供生产、消费的标准接口。
  3. 消费者宿主(ConsumerHost):管理消费者线程的生命周期,处理消息分发。
  4. 生命周期控制(CancellationToken):支持优雅关闭。

消息模型的设计决定了队列的通用性。不要把所有消息塞进一个队列里,否则不同类型的消息混在一起,消费者要做大量类型判断,耦合度很高。更合理的做法是按消息类型拆分队列,或者在消息体中携带 Topic/EventType 字段,由分发器做路由。

我这边的实践是做一个轻量级的发布订阅模型,每个消息包含一个事件类型,队列管理器内部维护一个"类型 -> 队列"的字典,生产者按类型投递,消费者按类型订阅。

public class Message { public string MessageId { get; set; } = Guid.NewGuid().ToString("N"); public string EventType { get; set; } public object Payload { get; set; } public DateTime Timestamp { get; set; } = DateTime.UtcNow; }

这里加MessageId是很重要的一步。虽然进程内通信不像分布式 MQ 那样有网络重投风险,但在消费者处理失败重试时,MessageId就是幂等判断的依据。这算是做消息队列的一个基本素养——不管规模大小,先埋好消息的唯一标识。

3.2 核心实现代码

队列管理器的核心代码没有多复杂,重点是几个设计的取舍。直接上代码。

using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Threading; using System.Threading.Tasks; /// <summary> /// 简易进程内消息队列 - 基于 ConcurrentQueue 实现 /// </summary> public class InMemoryMessageQueue : IDisposable { // 每个事件类型对应独立的队列 private readonly ConcurrentDictionary<string, ConcurrentQueue<Message>> _queues; // 每个队列对应的信号量,用于通知消费者有新消息 private readonly ConcurrentDictionary<string, SemaphoreSlim> _signals; // 维护当前所有消费者的取消令牌源 private readonly CancellationTokenSource _cts = new CancellationTokenSource(); // 记录是否有消费者在运行 private int _consumerCount; public InMemoryMessageQueue() { _queues = new ConcurrentDictionary<string, ConcurrentQueue<Message>>(); _signals = new ConcurrentDictionary<string, SemaphoreSlim>(); } /// <summary> /// 生产消息:投递到指定事件类型的队列 /// </summary> public void Publish(string eventType, object payload) { if (string.IsNullOrWhiteSpace(eventType)) throw new ArgumentException("事件类型不能为空", nameof(eventType)); var queue = _queues.GetOrAdd(eventType, _ => new ConcurrentQueue<Message>()); var signal = _signals.GetOrAdd(eventType, _ => new SemaphoreSlim(0, int.MaxValue)); var message = new Message { EventType = eventType, Payload = payload }; queue.Enqueue(message); signal.Release(); } /// <summary> /// 订阅消息:启动一个后台消费者线程,按事件类型消费 /// </summary> public void Subscribe(string eventType, Action<Message> handler, int consumers = 1) { if (handler == null) throw new ArgumentNullException(nameof(handler)); var queue = _queues.GetOrAdd(eventType, _ => new ConcurrentQueue<Message>()); var signal = _signals.GetOrAdd(eventType, _ => new SemaphoreSlim(0, int.MaxValue)); for (int i = 0; i < consumers; i++) { Interlocked.Increment(ref _consumerCount); var thread = new Thread(() => ConsumerLoop(eventType, queue, signal, handler, _cts.Token)) { IsBackground = true, Name = $"MQ-Consumer-{eventType}-{i}" }; thread.Start(); } } private void ConsumerLoop( string eventType, ConcurrentQueue<Message> queue, SemaphoreSlim signal, Action<Message> handler, CancellationToken token) { while (!token.IsCancellationRequested) { // 有消息时立即处理;没有消息时阻塞等待信号,避免 CPU 空转 if (queue.TryDequeue(out var message)) { try { handler(message); } catch (Exception ex) { // 异常处理很关键:默认记录日志后继续,不要让单个消息的失败拖垮消费者线程 Console.WriteLine($"[MQ] 消费消息失败: {ex}"); } continue; } // 队列为空,等待生产者信号,最多等待 500ms,定期检查取消请求 try { signal.Wait(token); } catch (OperationCanceledException) { break; } catch (Exception) { break; } } Interlocked.Decrement(ref _consumerCount); } /// <summary> /// 获取某类型队列的当前积压数量,用于监控 /// </summary> public int GetQueueLength(string eventType) { if (_queues.TryGetValue(eventType, out var queue)) return queue.Count; return 0; } public void Dispose() { _cts.Cancel(); // 等待消费者线程退出 while (Volatile.Read(ref _consumerCount) > 0) Thread.Sleep(20); _cts.Dispose(); } }

代码不长,也就 150 行左右,但每一块都有讲究。

ConcurrentDictionary按事件类型维护多个独立队列的好处是:不同类型的消息互不干扰,某一类消息的生产速度快不会挤占其他类型的资源。SemaphoreSlim(0, int.MaxValue)是这里的点睛之笔——它的初始计数是 0,生产者每次Release()就相当于投递了一个"有消息"的信号;消费者Wait()到信号后去队列里取消息。这样消费者线程在无消息时真正处于睡眠状态,完全不消耗 CPU。

3.3 生产者、消费者接入示例

代码写完,上实际使用的例子。模拟上一节说的上位机串口采集场景。

// 初始化消息队列 var mq = new InMemoryMessageQueue(); // 订阅:UI 刷新线程,单消费者 mq.Subscribe("DeviceData", msg => { var data = (DeviceData)msg.Payload; Console.WriteLine($"[UI] 刷新界面,显示设备温度 {data.Temperature:F2}°C"); }, consumers: 1); // 订阅:告警检测线程,双消费者提高吞吐 mq.Subscribe("DeviceData", msg => { var data = (DeviceData)msg.Payload; if (data.Temperature > 80) Console.WriteLine($"[WARN] 设备 {data.DeviceId} 温度过高 {data.Temperature:F2}°C"); }, consumers: 2); // 模拟串口采集线程 var cts = new CancellationTokenSource(); var collectThread = new Thread(() => { var random = new Random(); while (!cts.IsCancellationRequested) { var data = new DeviceData { DeviceId = "DEV-001", Temperature = 60 + random.NextDouble() * 30 }; mq.Publish("DeviceData", data); Thread.Sleep(50); // 模拟串口采集间隔 50ms } }); collectThread.IsBackground = true; collectThread.Start(); Console.WriteLine("消息队列运行中,按回车退出..."); Console.ReadLine(); cts.Cancel(); mq.Dispose();

这里有一个非常实用的设计——同一个事件类型可以被多个维度订阅。UI 刷新和告警检测都关注DeviceData,它们彼此独立,各自维护自己的消费者线程,如果有一天要增加新的处理模块,只需要再调一次Subscribe,完全不需要改动生产者的代码。这就是消息队列对抗业务变更的典型优势:发布者和订阅者完全解耦。

实测下来,这个实现单队列跨线程吞吐量能做到百万级消息/秒以上,对进程内通信来说是绰绰有余的。

4. 实战中的常见问题与排障思路

4.1 消息丢失、重复消费与顺序性

这三个问题是消息队列里被问得最多的,进程内队列也一样躲不开。

消息丢失:在进程内队列,"丢失"主要发生在两个地方。第一个是应用程序崩溃,内存队列里的所有消息都没了——这是内存队列的天然属性,解决思路是业务上接受"至少一次投递"的语义,或者对关键消息做好持久化后再发。第二个是消费者处理抛出异常,我在ConsumerLoop里默认是catch后直接跳过,这其实是主动丢弃问题消息的姿势。如果业务要求不能丢,应该在catch里做重试或者转移到"死信队列"。我在生产项目里是这么处理的:先判断消息是否已重试超过 3 次,没超就重新入队,超了就记录日志并丢弃。

catch (Exception ex) { if (message.RetryCount < 3) { message.RetryCount++; queue.Enqueue(message); signal.Release(); // 重新投递,让其他消费者或本消费者继续处理 Console.WriteLine($"[MQ] 消息重试第 {message.RetryCount} 次: {ex.Message}"); } else { Console.WriteLine($"[MQ] 消息处理失败已丢弃: {message.MessageId}"); } }

重复消费:进程内队列如果只有一个消费者,出队即移除,通常不会重复消费。但引入重试机制后,消息重新入队就可能被第二个消费者处理,造成重复。解决方案有两种:一是给消费者线程加处理锁,确保同一时刻只有一个线程在处理某条消息的重试;二是利用MessageId做幂等,消费者在处理前先记录或校验MessageId,已经处理过的就直接跳过。第二种方案更通用,推荐优先做。

消息顺序性:ConcurrentQueue<T>本身是 FIFO 的,保证入队顺序。但如果你启动了多个消费者线程,消息的分发顺序就无法保证了——线程 1 拿到消息 1 还在处理,线程 2 可能已经处理完消息 2。严格按顺序处理的办法是把consumers参数设为 1,单线程消费。但单线程消费会降低吞吐量,所以实际项目中要在"有序"和"高效"之间做权衡。我的做法是:按消息业务键做哈希分片,同一个设备的数据路由到同一个消费者线程,这样既保证了每个设备的消息有序,又让多个设备的处理并行。

4.2 内存膨胀与消费积压

进程内队列最容易被忽视的是内存监控。任务系统高峰期,生产者生产速度远大于消费者处理速度,队列里的消息数量就会不断增长,占用内存。生产环境里我见过因为积压消息太多导致进程内存飙升到几个 GB 的例子。

解决办法是在Publish方法里加积压保护。就像下水道的溢流阀一样,当消息积压超过阈值时,采取降级策略。这里给出一个简单的限流逻辑:

public bool Publish(string eventType, object payload, int maxQueueSize = 10000) { var queue = _queues.GetOrAdd(eventType, _ => new ConcurrentQueue<Message>()); // 积压保护:超过阈值直接拒绝新消息,防止内存膨胀 if (queue.Count > maxQueueSize) return false; var signal = _signals.GetOrAdd(eventType, _ => new SemaphoreSlim(0, int.MaxValue)); queue.Enqueue(new Message { EventType = eventType, Payload = payload }); signal.Release(); return true; }

调用方根据返回值决定是否走降级逻辑,比如丢弃消息、写入日志、或者提示系统繁忙。

但注意,queue.Count是弱一致性的,高并发下可能略小于实际值。追求更精准的话,可以用Interlocked维护单独的计数器变量。我在项目里对精度要求较高时是这样做的:在Publish后执行Interlocked.Increment(ref _counter),消费者取走消息后执行Interlocked.Decrement(ref _counter)。

监控积压:在生产环境,强烈建议每隔一段时间输出队列长度变化。我当时做了一个简单的定时监控线程:

var monitor = new Thread(() => { while (!cts.IsCancellationRequested) { foreach (var (eventType, queue) in mq.GetQueueSnapshots()) { Console.WriteLine($"[MONITOR] 队列 {eventType} 积压 {queue.Count} 条"); } Thread.Sleep(5000); } });

别看这些监控代码不起眼,线上系统出问题时,它们往往是最先暴露问题的哨兵。

4.3 锁、死锁与阻塞陷阱

ConcurrentQueue<T>本身是无锁的,但实际使用中死锁问题可能出在我们自己的代码上。

最常见的坑是:消费者处理器内部又去调用了同一个消息队列的Publish。比如一个消息触发了业务逻辑,业务逻辑又发了一条新消息。如果消费者线程数是 1,并且新消息与当前消息是同一个事件类型,就会产生隐式的顺序依赖,但不会死锁。真正会死锁的场景是:消费者处理消息时等待另一个队列的信号,而另一个队列的消费者又在等待本队列的信号,形成循环等待。

解决办法是:不要让消费者处理器内部同步等待其他队列的消费结果。如果确实有跨队列依赖,应该用异步回调的方式,或者把关联消息合并到同一个队列中,从设计上消除循环依赖。

另外一个隐蔽的陷阱是SemaphoreSlim.Wait(token)的使用。SemaphoreSlim在Dispose之后调用Wait会抛ObjectDisposedException,所以优雅关闭时,顺序非常重要:先Cancel取消令牌,让消费者线程从Wait中退出,等所有线程完全退出后再Dispose信号量。我上面代码中的Dispose方法就是这个顺序:

public void Dispose() { _cts.Cancel(); // 1. 通知消费者退出 while (Volatile.Read(ref _consumerCount) > 0) Thread.Sleep(20); // 2. 等待所有消费者线程安全退出 _cts.Dispose(); // 3. 最后释放资源 }

这个等待消费者退出的过程也很有讲究。Thread.Sleep(20)是轮询等待,如果消费者数量多,全部退出需要一些时间。更精细的做法是使用ManualResetEventSlim来等待,但为了保持代码简洁,轮询也是一种可以接受的方案,前提是消费者线程的finally中确保_consumerCount一定会递减。如果消费者处理消息时陷入死循环,这里的关闭就会卡住,所以开发时也要给消费者处理器设置超时机制。


整套代码跑下来,你会发现进程内消息队列远比想象中简单,但也比想象中更考验细节——积压监控、消费者异常处理、幂等设计、优雅关闭,每一步都是生产环境的必修课。我最初接触ConcurrentQueue<T>时,也觉得它只是替代Queue<T>加锁的线程安全版本,直到把它放在真实项目里当消息队列用,才体会到并发编程那些"看不见的坑"才是真正值得花时间的部分。如果你也在做上位机、后台任务调度或者实时数据分发,完全可以照着这套架构改造一份自己的版本,遇到问题欢迎一起交流踩坑经验。

返回列表