做工控上位机的兄弟应该都有这种体验程序里同时要处理PLC采集、数据库写入、UI刷新、日志落盘几个模块互相调用业务代码越写越拧巴。后来我在一个扭矩采集项目里因为要在UI上实时显示测量值同时后台还要落库、做统计汇总直接把采集线程和各个消费方耦合在一起结果界面一卡一卡的数据库写入稍微慢一点整个采集链路都跟着堵。被逼着搞了一个基于ConcurrentQueue的简易进程内消息队列问题一下子清爽了。这篇文章就把这套方案的完整思路、核心代码和踩过的坑都写出来适合有一定C#基础、想在单进程里做异步解耦但又不想引入RabbitMQ这类重型中间件的朋友参考。1. 为什么用ConcurrentQueue搭消息队列场景驱动1.1 没有消息队列时模块间通信有多痛很多C#新手习惯用最简单的方式写模块通信A模块直接调用B模块的方法。这种写法在业务简单时没毛病一旦模块多起来问题就非常明显。第一是耦合。采集模块要同时调UI刷新、数据库写入、日志记录那它就知道了所有下游模块的存在。今天加一个报表模块就要改采集模块明天去掉一个统计模块又要改。时间一长调用关系像蜘蛛网一样。第二是阻塞。方法调用是同步的采集线程要等数据库写入完成才能继续采集。数据库偶尔慢一下采集线程就被拖住后面的数据全积压在线程栈上UI自然卡顿。很多上位机程序越跑越卡不是机器不行是线程被同步调用卡死了。第三是失败传播。下游模块抛个异常直接炸穿上游调用链采集线程跟着挂掉。哪怕你写了try-catch也只是保证不崩溃可下游恢复后的数据补偿还是得自己手写。事件(event)能解决一部分耦合问题但实话说不够。event本身不支持消息堆积发布太快消费不过来就直接丢多个订阅者的处理顺序没法控制也没法做重试。你要实现“采集数据来了先放个地方谁需要谁自己来取取走之后慢慢处理处理失败还能重新来”这种机制就得自己建一个消息队列。1.2 进程内消息队列解决哪些实际问题进程内消息队列说白了就是在同一个进程里开辟一块“缓冲区”生产者只管往里扔消息消费者按自己的节奏取消息处理。拿我那个扭矩项目举例。拧紧枪通过PLC把测量值推给上位机每分钟好几千条数据。采集线程收到数据后只做一件事扔进队列。UI线程从队列里取数据刷新界面数据库线程取数据批量落库报表线程取数据做统计分析。三个消费者各取所需互不拖累。采集线程永远不会因为数据库慢而停下UI最多晚刷新几十毫秒但不会卡死。除了解耦它还能做到削峰。瞬时爆发的一万条消息可以先进队列消费者用稳定的速率慢慢处理避免瞬时压力打垮数据库连接或UI线程。别小看这个“进程内”三个字。跨进程的RabbitMQ、Kafka、Pulsar都很强大但部署、运维、网络开销都有成本。单体应用、上位机软件、桌面工具这类场景进程内队列完全够用还不用额外引入中间件部署包干净利落。1.3 为什么选ConcurrentQueue而不是ListlockC#里可选的无锁/线程安全集合不少但最适合当消息队列底座的确实是ConcurrentQueue。先看我在选型时对比过的几个方案方案线程安全阻塞等待内置背压控制实现复杂度适用场景ListT lock需要自己加锁无需要轮询需要自己实现较低简单缓存不推荐做队列BlockingCollectionT安全有Take阻塞有可限制容量低多数生产代码推荐直接用ConcurrentQueueT安全无TryDequeue立即返回需要自己实现中学习底层原理、需要完全控制逻辑时ChannelT安全有ReadAsync阻塞有Bounded限制低新项目首选.NET官方推荐ConcurrentQueue是微软提供的无锁并发队列入队出队都是原子的多个线程同时操作不会出问题。它的TryDequeue方法在线程安全和性能上做了很好的平衡但也有个特点队列空的时候它不会阻塞而是返回false。这个“不阻塞”在有些人眼里是缺点轮询会空转浪费CPU但在我看来恰好是优点——它把“等待”的控制权完全交给你你想让它等多久就等多久想用信号量唤醒还是用定时轮询你自己定。这也是为什么我选择ConcurrentQueue而不是直接用BlockingCollection目的不是造一个最省事的队列而是把队列的机制彻底搞清楚。当时我也认真考虑过ChannelT它确实更现代写起来也更简洁官方推荐新项目使用。但我在这个项目里的目的是做一个完全能掌控的简易消息队列同时还想让自己把无锁队列、同步机制这些基础吃透。如果你赶工期直接用Channel或BlockingCollection肯定更快这篇文章的代码思路照样能迁移过去。2. ConcurrentQueue底层逻辑明白它才能用好它2.1 无锁并发是怎么做到的很多人把ConcurrentQueue当黑盒用入队出队完事。它内部其实用的是CASCompare-And-Swap比较并交换这套无锁并发技术理解这一点对用好它非常重要。CAS的原理可以类比成高铁站的取号机多个窗口的乘客同时按取号按钮机器内部通过硬件指令保证只有一个请求能被接受其他人拿到的号各不相同。ConcurrentQueue内部维护着头指针和尾指针入队时用Interlocked.CompareExchange这类原子操作更新指针不需要像ListT那样用一个全局lock锁住整个集合。锁竞争的代价在线程一多时非常明显而CAS失败就重试轻量得多。存储结构上ConcurrentQueue不是一整块连续内存而是由若干个小段Segment组成的链表。每个段里装固定数量的元素段满了就新开一段段空了在一定条件下可回收。这样的好处是入队和出队操作只需锁定头尾附近的一小段区域并发冲突概率大大降低。这也是它比用一个List加一把大锁要快的根本原因。理解这层之后你就知道为什么ConcurrentQueue的很多用法和普通集合不一样了。它不是那种“拿索引随便读”的结构你只能从队尾入队、从队头出队或者看看队头是什么想要随机访问中间某个元素是不行的。2.2 线程安全边界与性能特点ConcurrentQueue给我们的安全保证很明确任意线程同时调用Enqueue、TryDequeue、TryPeek都没问题每个操作都是原子的。但你得知道它的边界在哪里。第一Count不是O(1)操作。它内部要遍历各段累加数量平时用一两次没事但如果你的代码在循环里高频调用Count来判断队列是否为空或长度是否超限性能会明显打折扣。我自己就优化过这种代码把Count 1000这种写法改成用信号量或计数器维护待处理数量效果立竿见影。第二ToArray之类的方法会把整个队列拷贝一份出来内存开销大不适合在生产路径上高频调用。需要批量查看消息时用TryDequeue逐个取出来更合理。第三引用类型入队时存储的是引用不是深拷贝。如果生产者入队后继续修改那个对象消费者读到的可能就是修改后的数据。所以在消息模型设计上我倾向于让MessageT不可变或者至少保证入队后不再改动原始对象。2.3 为什么FIFO顺序适合作为消息队列底座ConcurrentQueue是严格FIFO的先入队的元素一定先出队。这个特性对消息队列极其重要因为很多业务场景要求消息处理有顺序。比如扭矩数据的累加统计你必须保证先采集到的数据先被处理否则统计结果就是乱的。但这里必须说清楚一个容易混淆的点入队顺序不等于消费顺序。当你有多个消费者线程同时TryDequeue时线程A可能抢走第一条线程B抢走第二条两条消息的最终处理顺序可能颠倒。如果你对顺序有严格要求就必须用单消费者模式或者给消息加序号在消费者侧做排序或校验。我在那个项目里UI刷新用的是单消费者顺序无关紧要的日志写入用了四个消费者算是按需取舍。另外还要理解FIFO和消息优先级的关系。ConcurrentQueue不支持优先级所有消息一视同仁。要做到高优先级先处理得自己维护多个队列或者干脆用多个Topic隔离后面我会专门讲。3. 消息队列的核心设计消息模型与整体架构3.1 消息模型不止是塞个对象刚开始我用最简单的做法队列里直接塞业务对象。比如扭矩值就放一个TorqueData类。结果一调试就发现麻烦出了问题不知道这条消息是谁产生的、什么时候产生的、处理失败重试了几次。所以后来我把消息统一包装成一个泛型MessageT代码如下public sealed record MessageT( string Id, string Topic, T Payload, DateTime Timestamp, int RetryCount 0);用record是因为它天然带值相等性和with表达式后面做重试时修改RetryCount非常方便。每个字段都有存在的必要Id是消息唯一标识用GUID生成。它最核心的作用是解决“重复消费”问题。消费者处理完一条消息后如果宕机或异常重试它有可能再次拿到同一条消息没有Id你就无法判断这是不是一条已经处理过的旧消息。Topic是主题名类似于“通道名”用来区分不同类型的业务消息。扭矩数据、日志数据、UI通知分别走不同的主题互不干扰。Timestamp是生产时间排查消息积压、计算消费延迟时全靠它。RetryCount是重试次数配合下一轮的处理失败重试策略使用。3.2 队列中枢多主题管理的架构设计单队列只能解决“一对多”的解耦但实际项目里往往是“多对多”采集模块发扭矩数据日志模块发运行日志UI模块发界面通知。如果所有消息都挤在一个队列里消费方取消息时还得逐个判断类型处理速度被拖慢还可能互相阻塞。所以我设计了一个队列中枢类ProcessMessageQueueT内部用ConcurrentDictionarystring, ...按主题管理多个独立队列每个主题都有自己独立的ConcurrentQueue和信号量。发布消息时指定Topic消费者订阅时也指定Topic互不干扰。这里有一个非常隐蔽的坑必须拿出来说ConcurrentDictionary.GetOrAdd看起来是线程安全的但在并发极高时同一个key可能被创建出多个value实例。如果你直接GetOrAdd(topic, () new QueueContext())两个线程同时进来可能一个拿到队列A另一个拿到队列B然后消息就分散到两条不同的队列上去了订阅方只监听其中一条消息就莫名其妙丢了。正确的做法是让value保存一个LazyQueueContext确保同一个key只会真正初始化一次private readonly ConcurrentDictionarystring, LazyQueueContext _topics new(); private QueueContext GetContext(string topic) { return _topics.GetOrAdd( topic, _ new LazyQueueContext(() new QueueContext()) ).Value; }Lazy的默认执行模式是线程安全的多个线程同时触发初始化时只有第一个线程会真正创建对象其余线程等待它完成。这样保证同一个主题永远对应同一个队列实例。3.3 消费者模型单消费者还是多消费者消费者模型在设计时就要想清楚因为不同模型对吞吐量和顺序性的影响是相反的。单消费者模型就是一个主题只有一个后台任务在不断TryDequeue消息。优点是严格保证FIFO顺序前一条没处理完后一条绝对不会开始缺点是吞吐量受限于单个处理任务的速率。适合UI刷新、状态流转这种对顺序敏感的场景。多消费者模型就是一个主题启动多个后台任务同时取消息。吞吐量上去了但处理顺序无法保证。适合日志写入、数据落库这类不要求顺序、只求快的场景。我这里把消费者的启动和管理单独封装成一个QueueConsumerHostT它接收一个处理委托和消费者数量内部启动若干后台任务循环处理消息。这样使用方不需要自己管理线程生命周期。核心设计思路是队列只管存取消费宿主管业务处理各司其职。另外必须设计好优雅退出机制。我用了CancellationTokenSource释放资源时先Cancel()再Task.WaitAll等待当前正在处理的消息完成。这样程序退出时能尽量把正在处理的消息处理完不会直接半路杀死。4. 逐步实现一个可运行的简易消息队列4.1 队列上下文队列加信号量的组合先看队列的中枢类内部的队列上下文结构private sealed class QueueContext { public ConcurrentQueueMessageT Queue { get; } new(); public SemaphoreSlim Signal { get; } new(0); public long PendingCount; }Signal是我解决“空轮询”问题的关键。前面说过ConcurrentQueue的TryDequeue在空队列时不会阻塞直接返回false。如果消费者拿不到消息就死循环重试CPU会空转如果Thread.Sleep一下又会让消息处理延迟几十毫秒。SemaphoreSlim在这里扮演“门铃”的角色。生产者每发布一条消息就按一下门铃Release消费者在队列为空时去睡觉WaitAsync门铃一响立刻醒来取消息。这样既有实时性又不浪费CPU。PendingCount是我用Interlocked维护的待处理消息计数专门用来判断队列积压程度。不用ConcurrentQueue.Count是因为前面说过它要遍历段性能不划算。每发布一条消息加一每成功出队一条减一通过Interlocked.Read读取在高频路径上比Count快得多。4.2 发布与消费的核心代码下面是队列中枢的完整实现public sealed class ProcessMessageQueueT { private readonly ConcurrentDictionarystring, LazyQueueContext _topics new(); public void Publish(string topic, T payload) { var context GetContext(topic); var message new MessageT( Id: Guid.NewGuid().ToString(N), Topic: topic, Payload: payload, Timestamp: DateTime.UtcNow); context.Queue.Enqueue(message); Interlocked.Increment(ref context.PendingCount); context.Signal.Release(); } public async TaskMessageT? ConsumeAsync(string topic, CancellationToken cancellationToken) { var context GetContext(topic); while (!cancellationToken.IsCancellationRequested) { if (context.Queue.TryDequeue(out var message)) { Interlocked.Decrement(ref context.PendingCount); return message; } try { await context.Signal.WaitAsync(cancellationToken); } catch (OperationCanceledException) { return null; } } return null; } private QueueContext GetContext(string topic) { return _topics.GetOrAdd( topic, _ new LazyQueueContext(() new QueueContext()) ).Value; } }Publish方法里Signal.Release()和PendingCount自增的顺序有讲究。先入队、再计数、最后释放信号量。如果先释放信号量消费者醒过来去取消息时队列可能还没入队成功虽然这种情况概率极低但逻辑上不够严谨。ConsumeAsync方法里有个小循环先尝试TryDequeue不管有没有取到拿到就返回没拿到就等信号量。为什么要先试一次再等因为信号量和队列不是完全原子统一的可能会出现“信号量允许了但队列已经被别的消费者抢空”的情况这时候多走一轮循环就能修正不会死等。取消机制也在这里处理好了。CancellationToken取消时WaitAsync会抛OperationCanceledException我捕获它并返回null消费者就能安全退出。同时在等待期间如果有人取消信号量计数不会被错误消耗这是SemaphoreSlim的语义保证的。4.3 消费宿主把线程生命周期管起来有了队列中枢还得有消费端的管理类不然每个业务方都得自己写while循环代码会散落各处。我设计了QueueConsumerHostTpublic sealed class QueueConsumerHostT : IDisposable { private readonly ListTask _workers new(); private readonly CancellationTokenSource _cts new(); public QueueConsumerHost( ProcessMessageQueueT queue, string topic, FuncMessageT, Task handler, int workerCount 1) { for (var i 0; i workerCount; i) { _workers.Add(Task.Run(() ConsumeLoopAsync(queue, topic, handler, _cts.Token))); } } private async Task ConsumeLoopAsync( ProcessMessageQueueT queue, string topic, FuncMessageT, Task handler, CancellationToken cancellationToken) { while (!cancellationToken.IsCancellationRequested) { var message await queue.ConsumeAsync(topic, cancellationToken); if (message is null) break; try { await handler(message); } catch (Exception ex) { // 记录日志并按重试策略决定是否重新入队 Console.WriteLine($[消费失败] {message.Id}: {ex.Message}); } } } public void Dispose() { _cts.Cancel(); Task.WaitAll(_workers.ToArray()); _cts.Dispose(); } }消费者数量我做成参数方便调用方按业务场景调整。Task.Run启动的后台任务在线程池上运行不会单独占一个专用线程资源利用率更好。Dispose方法里先Cancel再WaitAll能保证正在执行中的handler消息处理完成。如果handler内部调用了Task.Delay之类最长等待时间就是业务处理的耗时一般可以接受。这套逻辑虽然简单但在我的上位机项目里做退出时很少丢消息。4.4 失败重试和死信队列业务处理总会有失败的时候。第一次写的时候我在catch里直接把异常吞了结果后来排查问题发现数据缺失都不知道丢在哪。我的处理策略是捕获异常后判断消息的RetryCount少于指定次数就重新入队超过就丢进一个“死信队列”并记录日志。借助之前record的with表达式重试逻辑写起来很清爽private const int MaxRetryCount 3; // 在消费宿主的catch块中调用队列的重试方法 public void RetryOrDeadLetter(ProcessMessageQueueT queue, MessageT message) { if (message.RetryCount MaxRetryCount) { queue.Reject(message.Topic, message); } else { Console.WriteLine($[死信] {message.Id} 重试超过 {MaxRetryCount} 次放弃处理); } }对应队列中枢里的Reject方法public void Reject(string topic, MessageT message) { var retried message with { RetryCount message.RetryCount 1 }; var context GetContext(topic); context.Queue.Enqueue(retried); Interlocked.Increment(ref context.PendingCount); context.Signal.Release(); }这里没有重新生成Id因为Id是用来幂等去重的。同一条消息重试多次Id必须保持一致消费者侧才能认出它来。而RetryCount每次加一用于控制最多重试次数避免坏消息在队列里无限循环。实际操作中我还会把死信消息序列化到本地文件或数据库表方便事后分析。简易队列能做到底层机制已经不错死信持久化不用太复杂记录关键字段即可。4.5 模拟压测一把真实场景代码写完之后我在本地跑了一个模拟场景一个生产者线程发布10万条扭矩测量数据三个消费者分别模拟UI刷新无耗时、数据库写入2毫秒耗时、统计计算1毫秒耗时。核心代码大致是这样var queue new ProcessMessageQueueTorqueData(); using var uiConsumer new QueueConsumerHostTorqueData( queue, torque, async msg { // 模拟UI刷新 await Task.CompletedTask; }, workerCount: 1); using var dbConsumer new QueueConsumerHostTorqueData( queue, torque, async msg { await Task.Delay(2); // 模拟数据库写入 }, workerCount: 3); using var statsConsumer new QueueConsumerHostTorqueData( queue, torque, async msg { await Task.Delay(1); // 模拟统计计算 }, workerCount: 2); for (var i 0; i 100_000; i) { queue.Publish(torque, new TorqueData(i, DateTime.UtcNow)); }注意这里三个消费者订阅的是同一个主题torque他们是竞争关系每条消息只会被其中一个消费者取走不是复制给三个消费者。如果需要“一条消息同时给三个消费者处理”那要用三个不同主题各发布一次或者设计广播机制。这也是进程内消息队列和事件总线的一个区别使用时要分清楚。压测结果很直观队列本身的入队出队完全不是瓶颈真正耗时都在消费者处理逻辑上。数据库写入2毫秒的延迟通过3个消费者并行吞吐比单消费者翻了接近三倍。UI刷新消费者因为消息到达及时界面图表刷新非常顺滑采集线程从头到尾没被阻塞过。5. 并发问题的避坑手册5.1 消息重复消费怎么破重复消费是消息队列绕不开的话题。我的场景里问题出现在消费者处理成功、但在返回消息确认之前进程被回收或者处理超时触发了重试。这时队列里已经拿不到那条消息了但业务系统可能已经执行了两次写入。解决重复消费的标准思路是“幂等”。简单说就是对同一条消息执行多次效果和只执行一次相同。实现方式有几种唯一索引去重数据库表里给消息Id建唯一索引重复插入会失败不会产生脏数据。状态标记记录已处理的消息Id集合处理前先查一下。进程内队列可以用ConcurrentDictionarystring, bool保存近期处理过的消息Id注意设置过期策略防止内存无限增长。业务幂等比如“把值设为X”这种操作天然幂等重复设置结果一样不需要额外处理。我在扭矩项目里给数据落库表加了消息Id字段和唯一索引重复消费根本写不进去问题直接根治。这个经验对所有消息队列方案都适用。5.2 消息丢失的三种姿势进程内队列最核心的局限就是消息保存在内存里进程一挂东西全没。这是没办法根治的但能尽量把丢失窗口缩小。第一种丢失姿势是消费者线程崩溃。如果ConsumeLoopAsync里的while循环没有包try-catch一次未捕获异常会让整个消费任务终止之后到达的消息全部积压没人处理。我踩过这个坑后来给循环体加了兜底catch确保单个消息失败不会导致消费者线程消失。第二种丢失姿势是业务处理异常后直接丢弃。你只写了个catch把异常打印出来但消息已经从队列里出队了本地日志再详细数据也回不来了。正确做法是走上面的重试和死信流程。第三种丢失姿势是退出时不等待消费任务完成。如果你直接Environment.Exit或主线程退出前不处理消费者任务正在队列里的消息和正在处理的消息就全丢了。所以我在消费宿主里实现Dispose时先Cancel再WaitAll强制等待正在处理的消息完成尽量降低丢失风险。要明白进程内队列不保证投递可靠性这是它的边界。需要严格不丢消息的场景必须上持久化的消息中间件。5.3 死锁没那么容易但小心这种写法自己写的队列代码很多人担心死锁。其实ConcurrentQueue本身是无锁的SemaphoreSlim也不容易出现死锁但业务代码里如果乱写还是能被自己坑到。最容易出问题的写法是在消费回调里等待另一个队列的消息。比如主题A的消费者在处理时同步等待从主题B取一条消息而主题B的消费者又反过来等待主题A这就形成了循环等待。虽然因为队列不会阻塞释放两个消费者互等的情况很少但多线程交叉等待一旦遇上信号量计数异常就真的会卡死。我给自己定了个规矩消费回调里不要做任何跨队列等待需要其他数据就去查数据库或缓存不要指望另一个队列配合你。消息队列是用来解耦的不是用来搞联动编排的。另外还有个和“死锁”同性质的坑消费者数量为0。如果你只发布消息没有启动任何QueueConsumerHost订阅对应主题消息就会一直积压。这在开发时不算问题上线时如果订阅逻辑被条件分支跳过就会很诡异。我后来在发布时增加了一个“无消费者”判断虽然不能阻止消息积压但至少能在日志里预警。5.4 背压队列积压了怎么办消息生产速度大于消费速度时队列会持续积压内存占用不断上涨。这就是背压问题。我的处理方式有两层。第一层是做好监控定期读取PendingCount超过阈值就记录预警日志。第二层是控制生产速度也就是在发布时做一个上限判断队列积压超过一定量就返回失败或丢弃新消息public bool TryPublish(string topic, T payload, int maxPendingCount 10000) { var context GetContext(topic); var current Interlocked.Read(ref context.PendingCount); if (current maxPendingCount) { Console.WriteLine($[背压] 主题 {topic} 积压 {current} 条拒绝新消息); return false; } Publish(topic, payload); return true; }用Interlocked.Read读计数而不是Count就是为了在高频调用下不拖累性能。maxPendingCount具体设多少要看业务能容忍多大的延迟。比如数据库消费者每秒处理500条积压10000条意味着最多晚20秒入库如果你的实时报表要求5秒内看到数据那10000就太大了。丢弃新消息这种策略要非常谨慎最好只用于非核心业务。核心数据宁可降低消费者处理速度也不要直接丢。5.5 什么时候别再用它这套进程内消息队列好归好但边界一定要清楚。第一需要跨进程通信时不要用它。多个独立进程之间要传数据ConcurrentQueue帮不上忙老老实实上RabbitMQ、Kafka、Pulsar或者Redis Stream。我这里说的是单进程内的模块解耦。第二需要消息持久化时不要用它。进程一退出内存队列全部清空。如果这些消息丢了会有经济损失或安全风险就必须用落盘方案。第三需要消息可靠投递时不要用它。进程内队列没法做消费确认、未消费消息转移这些机制。你需要的是专业的消息中间件而不是自己造的轮子。第四消息量极大且消费方需要复杂路由、延时、优先级时也别硬造。ConcurrentQueue能做FIFO做不了优先级队列做不了延时消息这些都要自己调度复杂度会迅速失控。我在实际项目里用它主要就是上位机和桌面工具这种单进程场景。一旦我发现项目开始需要多个服务节点协同或者有人提出“服务重启后消息不能丢”我就会第一时间换方案绝不硬扛。6. 常见问题速查表与个人体会6.1 问题排查速查表现象可能原因检查方法解决办法队列积压暴涨消费者线程太少或处理耗时太大查看PendingCount和每个handler的执行耗时增加workerCount优化消费者内部逻辑TryPublish加背压保护CPU飙高消费循环变成了空轮询检查代码是否用了while(true){TryDequeue}无等待用SemaphoreSlim.WaitAsync替代自旋或改用BlockingCollection消息处理重复重试机制导致同一条消息被消费多次查看消息Id在数据库记录中的重复次数消费侧做Id幂等数据库建唯一索引消费者退出后消息不再处理消费任务因未捕获异常退出查看进程日志有无异常堆栈在ConsumeLoopAsync里加catch保护检查while循环外部的异常消息偶尔消失业务处理失败后直接丢弃检查catch块里是否只写了日志走Reject重试流程超过次数进死信队列UI线程卡顿把UI刷新逻辑塞到了消费者高并发回调里检查handler是否用了Control.Invoke又同时耗时长消费线程只负责取消息UI刷新通过Dispatcher.InvokeAsync异步调度6.2 几个真实项目里的小技巧用久了你就会积累一些文档里不写的细节。我分享几个最实用的。第一任何消息处理链路都要有日志追踪。我在MessageT里放了Id后日志格式统一成“消息Id 主题 阶段”排查问题直接按Id搜整个处理链路一目了然。没有这个线上出现积压和重复消费时你根本不知道是哪条消息出了问题。第二队列健康状态要可视化。哪怕只是控制台输出或写一个状态文件定时把每个主题的待处理量、最近消费时间、最大延迟打出来也比出问题后盲猜强得多。我的项目里做了一个一分钟打印一次队列状态的计时器靠它发现过好几次消费者悄悄退出的问题。第三批量处理是提升吞吐的利器。数据库写入场景下一条一条INSERT效率很低。我改成消费者攒够100条或超过500毫秒就批量写入一次吞吐直接翻了几倍。ConcurrentQueue的出队逻辑支持这种模式循环TryDequeue取出一批再统一处理。第四消息内容尽量小而美。队列里塞大对象会占内存还会拉高GC压力。能放Id就放Id需要完整数据可以消费时再查。我用扭矩数据做消息时只放测量值和关键字段不做复杂对象嵌套GC明显稳定。6.3 最后说点个人体会这套基于ConcurrentQueue的进程内消息队列我用了好几个项目边用边改到现在算是稳定了。它不华丽功能也很基础但恰恰因为简单我才能完全掌握它的每一处行为遇到问题能一眼定位。对比直接用Channel自己搭这套东西多写了不少代码但收获是实打实的我对无锁并发、信号量同步、消费者生命周期这些东西的理解深了一大截。有一次面试聊消息队列重复消费怎么解决我把Message.Id幂等方案讲了一遍面试官还追问了信号量和队列事务之间的匹配问题我正好用这套实现里的经验对答如流。如果你也在做上位机、桌面工具或者单体服务又被模块通信搞得很烦可以花一晚上把这里面的代码敲一遍改成你自己业务的Topic和消息体。用上之后你会明显感觉到代码清爽了出问题也好查了。真到了需要跨进程、持久化的那天你带着这套理解再去学RabbitMQ或Kafka也会轻松很多。
阅读完成 · 觉得有帮助?