.NET ConcurrentDictionary在RabbitMQ工作线程中的状态一致性问题
问题根源与解决方案
核心问题
关闭窗体后,旧窗体实例的RabbitMQ消费者并未被销毁,仍在后台监听队列。重新打开窗体后,队列存在多个消费者(旧窗体+新窗体实例),导致后续消息被分散消费:第二条消息被旧窗体实例接收并添加到它的_DocumentsList(你无法观测到这个旧字典),而新窗体的字典只收到第一条,因此出现"第二条消息添加后计数仍为1"的假象。
关键原因分析
ReceiveAsync调用BasicConsume后,RabbitMQ通道会持有消费者引用,即使窗体关闭,只要通道未关闭、消费者未被取消,旧实例的消费者会持续监听队列。- 异步void方法(
Async Sub ProcessQueueMessage)存在异常捕获困难、线程上下文管理混乱的问题,进一步加剧了实例生命周期的不可控性。 - 未在窗体关闭时主动取消消费者注册,导致旧实例无法被GC回收,一直后台运行。
分步解决方案
1. 为消费者添加取消机制
修改第三方库的ReceiveAsync方法,返回消费者标签(consumerTag),并新增取消消费的方法:
// 改造ReceiveAsync,返回consumerTag用于后续取消 public async Task<string> ReceiveAsync<T>(string queue, Func<T, Task> onMessage) { _channel.QueueDeclare(queue, true, false, false); var consumer = new AsyncEventingBasicConsumer(_channel); consumer.Received += async (s, e) => { try { byte[] array = new byte[e.Body.Span.Length]; e.Body.Span.CopyTo(array); var jsonSpecified = Encoding.UTF8.GetString(array); var item = JsonConvert.DeserializeObject<T>(jsonSpecified); // 直接调用异步方法,避免不必要的Task.Run线程切换 await onMessage(item); } catch (Exception ex) { // 记录日志而非抛出,避免影响其他消息处理 Debug.WriteLine($"消息处理失败: {ex.Message}"); } }; // 注册消费者并返回consumerTag var consumerTag = _channel.BasicConsume(queue, true, consumer); return consumerTag; } // 新增取消消费的方法 public void CancelConsume(string consumerTag) { if (_channel?.IsOpen == true && !string.IsNullOrEmpty(consumerTag)) { _channel.BasicCancel(consumerTag); } }
2. 改造窗体代码,在关闭时取消消费者
将异步void方法改为异步函数,并在窗体关闭时主动取消消费:
Private Property _DocumentsList As New ConcurrentDictionary(Of Double, Doc) Private Property _API As API Private _consumerTag As String ' 保存消费者标签 Public Sub New(ByVal inAPI As API) InitializeComponent() _API = inAPI RegisterQueueHandler() End Sub Private Async Sub RegisterQueueHandler() Dim queueName = "ProcessorQueue" & _API._EnvironmentID.ToString ' 接收返回的consumerTag _consumerTag = Await _API._BusQueueControl.ReceiveAsync(Of ProcessorQueue)(queueName, AddressOf ProcessQueueMessage) End Sub ' 将Async Sub改为Async Function,避免异步void的问题 Private Async Function ProcessQueueMessage(q As ProcessorQueue) As Task Dim processorService As New ProcessorService(_API) Dim newDoc As Doc = Await processorService.MergeAndQueueForProcessing(q) Dim success As Boolean = _DocumentsList.TryAdd(newDoc.Number, newDoc) If success Then Debug.WriteLine("Added Document Number " & newDoc.Number.ToString) Else Debug.WriteLine("Failed to add " & newDoc.Number.ToString()) End If End Function ' 窗体关闭时取消消费者注册 Private Sub frmDocSplitting_FormClosing(sender As Object, e As FormClosingEventArgs) Handles MyBase.FormClosing _API._BusQueueControl.CancelConsume(_consumerTag) End Sub
3. 额外优化点
- 避免使用
Async Sub:异步void方法的异常无法被上层捕获,容易导致程序崩溃或隐藏错误,统一使用Async Function返回Task。 - 移除不必要的
Task.Run:ProcessQueueMessage本身是异步方法,无需再用Task.Run包装,减少线程切换开销。 - 检查通道生命周期:如果
_API._BusQueueControl的通道是全局共享的,确保取消单个消费者不会影响其他业务;如果是每个窗体实例独立通道,关闭窗体时可以同时关闭通道。
验证方案
- 打开窗体,发送2条消息,确认字典计数正常变为2。
- 关闭窗体,重新打开,再次发送2条消息:
- 第一条消息被新窗体接收,计数变为1。
- 第二条消息也被新窗体接收,计数变为2。
- 旧窗体实例的消费者已被取消,不会再消费消息。
内容的提问来源于stack exchange,提问作者Jurgen Cuschieri
相关产品推荐
相关产品推荐

