You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

.NET ConcurrentDictionary在RabbitMQ工作线程中的状态一致性问题

问题根源与解决方案

核心问题

关闭窗体后,旧窗体实例的RabbitMQ消费者并未被销毁,仍在后台监听队列。重新打开窗体后,队列存在多个消费者(旧窗体+新窗体实例),导致后续消息被分散消费:第二条消息被旧窗体实例接收并添加到它的_DocumentsList(你无法观测到这个旧字典),而新窗体的字典只收到第一条,因此出现"第二条消息添加后计数仍为1"的假象。

关键原因分析

  1. ReceiveAsync调用BasicConsume后,RabbitMQ通道会持有消费者引用,即使窗体关闭,只要通道未关闭、消费者未被取消,旧实例的消费者会持续监听队列。
  2. 异步void方法(Async Sub ProcessQueueMessage)存在异常捕获困难、线程上下文管理混乱的问题,进一步加剧了实例生命周期的不可控性。
  3. 未在窗体关闭时主动取消消费者注册,导致旧实例无法被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的通道是全局共享的,确保取消单个消费者不会影响其他业务;如果是每个窗体实例独立通道,关闭窗体时可以同时关闭通道。

验证方案

  1. 打开窗体,发送2条消息,确认字典计数正常变为2。
  2. 关闭窗体,重新打开,再次发送2条消息:
    • 第一条消息被新窗体接收,计数变为1。
    • 第二条消息也被新窗体接收,计数变为2。
    • 旧窗体实例的消费者已被取消,不会再消费消息。

内容的提问来源于stack exchange,提问作者Jurgen Cuschieri

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.17 10:29:52