BlockingCollection的GetConsumingEnumerable停止消费致集合膨胀问题排查
问题描述
在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

