如何实现支持总内存权重限制的线程安全WeighedBlockingCollection<T>集合?
如何实现支持总内存权重限制的线程安全WeighedBlockingCollection集合?
要实现一个符合需求的WeighedBlockingCollection<T>,我们可以基于.NET内置的BlockingCollection<T>(复用其成熟的多生产者/消费者、阻塞、FIFO等核心逻辑),结合手动维护的权重计数和同步原语来控制总内存占用,同时满足延迟创建对象的要求。以下是完整的线程安全实现,完全匹配你定义的接口和行为要求:
using System; using System.Collections.Concurrent; using System.Collections.Generic; using System.Threading; public class WeighedBlockingCollection<T> { private readonly long _maxTotalWeight; private long _currentTotalWeight; private readonly object _syncLock = new object(); private readonly BlockingCollection<(T Item, long Weight)> _innerCollection; private bool _isAddingCompleted; public WeighedBlockingCollection(long maximumTotalWeight) { if (maximumTotalWeight <= 0) throw new ArgumentOutOfRangeException( nameof(maximumTotalWeight), "Maximum total weight must be a positive value."); _maxTotalWeight = maximumTotalWeight; _currentTotalWeight = 0; _isAddingCompleted = false; // 用ConcurrentQueue作为底层存储,严格保证FIFO顺序 _innerCollection = new BlockingCollection<(T Item, long Weight)>(new ConcurrentQueue<(T Item, long Weight)>()); } public void Add(long itemWeight, Func<T> itemFactory) { // 先验证输入参数合法性 if (itemWeight < 0) throw new ArgumentOutOfRangeException( nameof(itemWeight), "Item weight cannot be negative."); if (itemWeight > _maxTotalWeight) throw new ArgumentOutOfRangeException( nameof(itemWeight), $"Item weight exceeds the maximum total weight of the collection ({_maxTotalWeight})."); if (itemFactory == null) throw new ArgumentNullException(nameof(itemFactory)); lock (_syncLock) { // 检查是否已经停止接受新元素 if (_isAddingCompleted) throw new InvalidOperationException("Adding has been completed for this collection."); // 循环等待直到有足够的权重空间(处理虚假唤醒) while (_currentTotalWeight + itemWeight > _maxTotalWeight) { if (_isAddingCompleted) throw new InvalidOperationException("Adding was completed while waiting for available space."); // 释放锁并进入等待队列,直到被消费者唤醒(按FIFO顺序) Monitor.Wait(_syncLock); } // 延迟创建对象:此时已确认有足够空间,由调用Add的生产者线程同步执行 T item = itemFactory(); // 将对象和对应权重加入内部阻塞集合 _innerCollection.Add((item, itemWeight)); // 更新当前总权重计数 _currentTotalWeight += itemWeight; } } public void CompleteAdding() { lock (_syncLock) { if (_isAddingCompleted) return; _isAddingCompleted = true; // 通知内部阻塞集合停止接受新元素 _innerCollection.CompleteAdding(); // 唤醒所有等待的生产者,让它们抛出异常终止等待 Monitor.PulseAll(_syncLock); } } public IEnumerable<T> GetConsumingEnumerable() { // 复用BlockingCollection的GetConsumingEnumerable逻辑:阻塞取元素,直到集合为空且添加完成 foreach (var (item, weight) in _innerCollection.GetConsumingEnumerable()) { // 释放当前元素占用的权重,并唤醒最早等待的生产者(保证FIFO顺序) lock (_syncLock) { _currentTotalWeight -= weight; // 唤醒一个等待的生产者(按调用顺序FIFO) Monitor.Pulse(_syncLock); } yield return item; } } }
核心设计细节与行为说明
权重控制与同步机制:
- 用
_currentTotalWeight(线程安全的long类型)维护当前集合的总权重,所有对该值的修改和读取都通过_syncLock锁保护,避免多线程竞争问题。 - 通过
Monitor.Wait和Monitor.Pulse实现生产者的阻塞等待和唤醒:生产者会等待直到有足够的权重空间,消费者取出元素后会唤醒最早等待的生产者(严格遵循FIFO顺序,不会优先处理小对象)。
- 用
延迟对象创建逻辑:
- 只有当确认集合有足够的权重空间时,才会调用
itemFactory创建对象,完全符合“延迟创建”的要求,且创建过程由调用Add的生产者线程同步执行,不会移交到消费者线程。
- 只有当确认集合有足够的权重空间时,才会调用
线程安全与多生产者/消费者支持:
- 所有共享状态(
_currentTotalWeight、_isAddingCompleted)都被lock保护,内部使用的BlockingCollection<(T, long)>本身就是为多生产者/消费者场景设计的线程安全集合。 - 支持任意数量的生产者和消费者,行为与内置
BlockingCollection<T>完全一致。
- 所有共享状态(
异常处理与边界情况:
- 当
itemWeight为负数、超过最大总权重时,会立即抛出ArgumentOutOfRangeException。 - 调用
CompleteAdding后,后续的Add调用会抛出InvalidOperationException,正在等待的生产者也会被唤醒并抛出异常。 - 若
itemFactory本身抛出异常,不会占用任何权重(因为对象创建失败后未加入集合,权重计数也未更新),不会影响其他生产者的正常操作。
- 当
FIFO顺序保证:
- 内部使用
ConcurrentQueue作为BlockingCollection的底层存储,保证元素的FIFO顺序;生产者的等待队列通过Monitor.Wait的内置FIFO机制实现,确保Add调用的顺序性,不会优先处理小对象。
- 内部使用
需求匹配验证
- ✅ 支持多生产者/消费者,完全线程安全
- ✅ 严格FIFO的
Add顺序,不会优先处理小对象 - ✅ 延迟对象创建,仅在有足够空间时由生产者线程同步执行
- ✅ 完全匹配
BlockingCollection<T>的核心行为(CompleteAdding、GetConsumingEnumerable的阻塞/结束逻辑) - ✅ 正确处理非法权重参数和
CompleteAdding后的操作 - ✅ 支持任意大的总权重(基于long类型,无int.MaxValue限制)
内容来源于stack exchange
相关产品推荐
相关产品推荐

