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

BlockingCollection的GetConsumingEnumerable停止消费致集合膨胀问题排查

BlockingCollection消费线程停止导致内存增长排查方案

问题描述

在Windows服务中使用BlockingCollection<string>,生产者每25ms向集合添加字符串,消费者线程通过GetConsumingEnumerable()循环消费,但该线程会突然停止消费,导致集合内存持续增长。经dotmemory分析确认内存占用集中在该集合,日志未记录线程停止原因,重启服务可临时解决但会复发。

提供代码

public delegate void EventHandler(string message);
 public event EventHandler OnEvent;

 private BlockingCollection<string> m_blockingCollection;

 // 每25ms调用一次
 private void Add(string msg)
 {
    m_blockingCollection.Add(msg);
 }
 // 服务启动时调用
 private void StartThread()
 {
    m_blockingCollection = new BlockingCollection<string>();
    System.Threading.Thread t = new System.Threading.Thread(ProducerConsumer)
    {
       IsBackground = true
    };
    t.Start();
 }
 private void ProducerConsumer()
 {
    try
      {
         foreach (var telegram in m_blockingCollection.GetConsumingEnumerable())
          {
             try
              {
                OnEvent?.Invoke(telegram);
              }
              catch (Exception ex)
              {
                 LogError(ex);
              }
          }
      }
      catch (Exception ex)
      {
         LogError(ex);
      }
      finally
      {
         LogWarning("Exiting producer consumer");
      }
 } 

排查与修复方案

1. 重点排查事件订阅者阻塞线程

消费者线程在调用OnEvent?.Invoke(telegram)时,若某个订阅者的处理逻辑长时间阻塞(比如同步调用外部服务超时、死锁、等待未释放的锁),会直接卡住整个消费循环,且不会抛出异常,因此日志无记录。

修复建议:

  • 将事件调用改为异步执行,避免阻塞消费线程:
    Task.Run(() => 
    {
        try { OnEvent?.Invoke(telegram); }
        catch (Exception ex) { LogError(ex); }
    }).ConfigureAwait(false);
    
  • 给每个事件订阅者的处理逻辑添加超时控制,防止无限阻塞。

2. 检查集合是否被意外标记为完成

若代码其他地方调用了m_blockingCollection.CompleteAdding(),GetConsumingEnumerable()会在消费完现有元素后退出循环,线程进入finally块输出退出日志。若你未看到该日志,可能是日志未正确捕获,或存在误调用CompleteAdding的代码。

排查建议:

  • 全局搜索代码,确认是否存在调用CompleteAdding的逻辑;
  • 在finally块的日志中补充线程ID、集合状态信息:
    LogWarning($"Exiting producer consumer, ThreadID: {Thread.CurrentThread.ManagedThreadId}, IsAddingCompleted: {m_blockingCollection.IsAddingCompleted}, Count: {m_blockingCollection.Count}");
    

3. 检查线程是否被意外终止

虽然线程设置为IsBackground = true,但如果进程中发生未处理异常、或其他线程调用了Thread.Abort(),会导致线程终止。ThreadAbortException在捕获后会自动重新抛出,可能导致线程直接退出而未被外层catch块完整记录。

修复建议:

  • 单独捕获ThreadAbortException并记录,同时调用Thread.ResetAbort()防止自动抛出:
    catch (ThreadAbortException ex)
    {
        LogError($"Consumer thread aborted: {ex.Message}", ex);
        Thread.ResetAbort();
    }
    
  • 在ProducerConsumer方法开头添加线程启动日志,记录线程ID,方便和退出日志对应。

4. 排查多线程初始化集合的问题

若StartThread被多次调用,会重新初始化m_blockingCollection,旧集合被丢弃,原消费线程继续消费旧集合,新元素却被加到新集合,导致新集合无人消费。

修复建议:

  • 给集合初始化加线程安全锁,确保只初始化一次:
    private readonly object _initLock = new object();
    private void StartThread()
    {
        lock (_initLock)
        {
            if (m_blockingCollection == null)
            {
                m_blockingCollection = new BlockingCollection<string>();
                Thread t = new Thread(ProducerConsumer) { IsBackground = true, Name = "BlockingCollectionConsumer" };
                t.Start();
            }
        }
    }
    

5. 增加消费过程日志监控

在消费循环内部添加日志,记录每次消费的元素和时间,确认是突然停止还是逐渐变慢,便于定位问题节点:

foreach (var telegram in m_blockingCollection.GetConsumingEnumerable())
{
    try
    {
        LogInfo($"Consuming message: {telegram}, ThreadID: {Thread.CurrentThread.ManagedThreadId}, Time: {DateTime.Now:yyyy-MM-dd HH:mm:ss.fff}");
        // 事件调用逻辑
    }
    // ...
}

内容的提问来源于stack exchange,提问作者Sai Kiran Vedire

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.03 04:25:38