【问题标题】:Creating a blocking Queue<T> in .NET?在 .NET 中创建阻塞队列<T>?
【发布时间】:2010-10-06 12:34:44
【问题描述】:

我有一个场景,我有多个线程添加到一个队列中,并且多个线程从同一个队列中读取。如果队列达到特定大小正在填充队列的所有线程将在添加时被阻塞,直到从队列中删除项目。

下面的解决方案是我现在正在使用的,我的问题是:如何改进?在我应该使用的 BCL 中是否存在已经启用此行为的对象?

internal class BlockingCollection<T> : CollectionBase, IEnumerable
{
    //todo: might be worth changing this into a proper QUEUE

    private AutoResetEvent _FullEvent = new AutoResetEvent(false);

    internal T this[int i]
    {
        get { return (T) List[i]; }
    }

    private int _MaxSize;
    internal int MaxSize
    {
        get { return _MaxSize; }
        set
        {
            _MaxSize = value;
            checkSize();
        }
    }

    internal BlockingCollection(int maxSize)
    {
        MaxSize = maxSize;
    }

    internal void Add(T item)
    {
        Trace.WriteLine(string.Format("BlockingCollection add waiting: {0}", Thread.CurrentThread.ManagedThreadId));

        _FullEvent.WaitOne();

        List.Add(item);

        Trace.WriteLine(string.Format("BlockingCollection item added: {0}", Thread.CurrentThread.ManagedThreadId));

        checkSize();
    }

    internal void Remove(T item)
    {
        lock (List)
        {
            List.Remove(item);
        }

        Trace.WriteLine(string.Format("BlockingCollection item removed: {0}", Thread.CurrentThread.ManagedThreadId));
    }

    protected override void OnRemoveComplete(int index, object value)
    {
        checkSize();
        base.OnRemoveComplete(index, value);
    }

    internal new IEnumerator GetEnumerator()
    {
        return List.GetEnumerator();
    }

    private void checkSize()
    {
        if (Count < MaxSize)
        {
            Trace.WriteLine(string.Format("BlockingCollection FullEvent set: {0}", Thread.CurrentThread.ManagedThreadId));
            _FullEvent.Set();
        }
        else
        {
            Trace.WriteLine(string.Format("BlockingCollection FullEvent reset: {0}", Thread.CurrentThread.ManagedThreadId));
            _FullEvent.Reset();
        }
    }
}

【问题讨论】:

  • .Net 如何有内置类来帮助解决这种情况。此处列出的大多数答案都已过时。请参阅底部的最新答案。查看线程安全的阻塞集合。答案可能已经过时,但它仍然是一个好问题!
  • 我认为学习 Monitor.Wait/Pulse/PulseAll 仍然是一个好主意,即使我们在 .NET 中有新的并发类。
  • 同意@thewpfguy。您需要了解幕后的基本锁定机制。另外值得注意的是,Systems.Collections.Concurrent 直到 2010 年 4 月才存在,然后仅在 Visual Studio 2010 及更高版本中存在。绝对不是 VS2008 支持的选项...
  • 如果您现在正在阅读本文,请查看 System.Threading.Channels 以了解针对 .NET Core 和 .NET 的多写入器/多读取器、有界、可选阻塞实现标准。

标签: c# .net multithreading collections queue


【解决方案1】:

如果您想要最大的吞吐量,允许多个读取器读取并且只有一个写入器写入,BCL 有一个称为 ReaderWriterLockSlim 的东西,它应该有助于精简您的代码...

【讨论】:

【解决方案2】:

这看起来很不安全(很少同步);像这样的东西怎么样:

class SizeQueue<T>
{
    private readonly Queue<T> queue = new Queue<T>();
    private readonly int maxSize;
    public SizeQueue(int maxSize) { this.maxSize = maxSize; }

    public void Enqueue(T item)
    {
        lock (queue)
        {
            while (queue.Count >= maxSize)
            {
                Monitor.Wait(queue);
            }
            queue.Enqueue(item);
            if (queue.Count == 1)
            {
                // wake up any blocked dequeue
                Monitor.PulseAll(queue);
            }
        }
    }
    public T Dequeue()
    {
        lock (queue)
        {
            while (queue.Count == 0)
            {
                Monitor.Wait(queue);
            }
            T item = queue.Dequeue();
            if (queue.Count == maxSize - 1)
            {
                // wake up any blocked enqueue
                Monitor.PulseAll(queue);
            }
            return item;
        }
    }
}

(编辑)

实际上,您需要一种关闭队列的方法,以便读者开始干净地退出 - 可能类似于 bool 标志 - 如果设置,则只会返回一个空队列(而不是阻塞):

bool closing;
public void Close()
{
    lock(queue)
    {
        closing = true;
        Monitor.PulseAll(queue);
    }
}
public bool TryDequeue(out T value)
{
    lock (queue)
    {
        while (queue.Count == 0)
        {
            if (closing)
            {
                value = default(T);
                return false;
            }
            Monitor.Wait(queue);
        }
        value = queue.Dequeue();
        if (queue.Count == maxSize - 1)
        {
            // wake up any blocked enqueue
            Monitor.PulseAll(queue);
        }
        return true;
    }
}

【讨论】:

  • 如何将等待更改为 WaitAny 并在构造时传入终止等待句柄...
  • @Marc- 如果您希望队列始终达到容量,则优化是将 maxSize 值传递给 Queue 的构造函数。你可以在你的类中添加另一个构造函数来适应它。
  • 为什么是 SizeQueue,为什么不是 FixedSizeQueue?
  • @Lasse - 它在Wait 期间释放锁,因此其他线程可以获取它。它在唤醒时回收锁。
  • 很好,正如我所说,有些东西我没有得到:) 这确实让我想重新审视我的一些线程代码......
【解决方案3】:

好吧,你可以看看System.Threading.Semaphore 类。除此之外 - 不,你必须自己做。 AFAIK 没有这样的内置集合。

【讨论】:

  • 我查看了它来限制访问资源的线程数,但它不允许您根据某些条件(如 Collection.Count)阻止对资源的所有访问。无论如何,AFAIK
  • 好吧,你自己做那部分,就像你现在做的那样。只需使用 Semaphore 代替 MaxSize 和 _FullEvent,您可以在构造函数中使用正确的计数对其进行初始化。然后,在每次添加/删除时调用 WaitForOne() 或 Release()。
  • 它和你现在拥有的没什么不同。更简单的恕我直言。
  • 你能给我一个例子来说明这个工作吗?我没有看到如何动态调整这种情况所需的信号量的大小。因为只有当队列已满时,您才能阻塞所有资源。
  • 啊,改变大小!为什么不马上说?好的,那么信号量不适合你。祝你好运!
【解决方案4】:

我还没有完全探索过TPL,但他们可能有适合您需求的东西,或者至少有一些 Reflector 素材可以从中获得一些灵感。

希望对您有所帮助。

【讨论】:

  • 我知道这是旧的,但我的评论是针对 SO 的新手,因为 OP 今天已经知道这一点。这不是答案,应该是评论。
【解决方案5】:

“如何改进?”

好吧,您需要查看类中的每个方法,并考虑如果另一个线程同时调用该方法或任何其他方法会发生什么。例如,您在 Remove 方法中放置了一个锁,但没有在 Add 方法中。如果一个线程在添加另一个线程的同时删除会发生什么? 坏事。

还要考虑一个方法可以返回第二个对象,该对象提供对第一个对象的内部数据的访问 - 例如,GetEnumerator。想象一个线程正在通过该枚举器,另一个线程同时正在修改列表。 不好。

一个好的经验法则是通过将类中的方法数量减少到绝对最小值来简化此操作。

特别是,不要继承另一个容器类,因为您将公开该类的所有方法,从而为调用者提供一种破坏内部数据或查看数据的部分完整更改的方法(同样糟糕,因为此时数据似乎已损坏)。隐藏所有详细信息,并对您如何允许访问它们完全无情。

我强烈建议您使用现成的解决方案 - 获取有关线程的书籍或使用 3rd 方库。否则,鉴于您正在尝试的内容,您将需要长时间调试代码。

另外,Remove 返回一个项目(例如,首先添加的项目,因为它是一个队列)而不是调用者选择特定项目不是更有意义吗?而当队列为空时,也许 Remove 也应该阻塞。

更新:Marc 的回答实际上实现了所有这些建议! :) 但我将把它留在这里,因为它可能有助于理解为什么他的版本如此改进。

【讨论】:

    【解决方案6】:

    这就是我为线程安全有界阻塞队列而来的操作。

    using System;
    using System.Collections.Generic;
    using System.Text;
    using System.Threading;
    
    public class BlockingBuffer<T>
    {
        private Object t_lock;
        private Semaphore sema_NotEmpty;
        private Semaphore sema_NotFull;
        private T[] buf;
    
        private int getFromIndex;
        private int putToIndex;
        private int size;
        private int numItems;
    
        public BlockingBuffer(int Capacity)
        {
            if (Capacity <= 0)
                throw new ArgumentOutOfRangeException("Capacity must be larger than 0");
    
            t_lock = new Object();
            buf = new T[Capacity];
            sema_NotEmpty = new Semaphore(0, Capacity);
            sema_NotFull = new Semaphore(Capacity, Capacity);
            getFromIndex = 0;
            putToIndex = 0;
            size = Capacity;
            numItems = 0;
        }
    
        public void put(T item)
        {
            sema_NotFull.WaitOne();
            lock (t_lock)
            {
                while (numItems == size)
                {
                    Monitor.Pulse(t_lock);
                    Monitor.Wait(t_lock);
                }
    
                buf[putToIndex++] = item;
    
                if (putToIndex == size)
                    putToIndex = 0;
    
                numItems++;
    
                Monitor.Pulse(t_lock);
    
            }
            sema_NotEmpty.Release();
    
    
        }
    
        public T take()
        {
            T item;
    
            sema_NotEmpty.WaitOne();
            lock (t_lock)
            {
    
                while (numItems == 0)
                {
                    Monitor.Pulse(t_lock);
                    Monitor.Wait(t_lock);
                }
    
                item = buf[getFromIndex++];
    
                if (getFromIndex == size)
                    getFromIndex = 0;
    
                numItems--;
    
                Monitor.Pulse(t_lock);
    
            }
            sema_NotFull.Release();
    
            return item;
        }
    }
    

    【讨论】:

    • 您能否提供一些代码示例,说明我如何使用该库对一些线程函数进行排队,包括如何实例化此类?
    • 这个问题/回答有点过时了。您应该查看 System.Collections.Concurrent 命名空间以获取阻塞队列支持。
    【解决方案7】:

    我刚刚使用响应式扩展解决了这个问题并记住了这个问题:

    public class BlockingQueue<T>
    {
        private readonly Subject<T> _queue;
        private readonly IEnumerator<T> _enumerator;
        private readonly object _sync = new object();
    
        public BlockingQueue()
        {
            _queue = new Subject<T>();
            _enumerator = _queue.GetEnumerator();
        }
    
        public void Enqueue(T item)
        {
            lock (_sync)
            {
                _queue.OnNext(item);
            }
        }
    
        public T Dequeue()
        {
            _enumerator.MoveNext();
            return _enumerator.Current;
        }
    }
    

    不一定完全安全,但非常简单。

    【讨论】:

    • 什么是主题?我的命名空间没有任何解析器。
    • 它是响应式扩展的一部分。
    • 不是答案。这根本无法回答问题。
    【解决方案8】:

    使用 .net 4 BlockingCollection,入队使用 Add(),出队使用 Take()。它在内部使用非阻塞 ConcurrentQueue。更多信息在这里Fast and Best Producer/consumer queue technique BlockingCollection vs concurrent Queue

    【讨论】:

      【解决方案9】:

      您可以在 System.Collections.Concurrent 命名空间中使用 BlockingCollection 和 ConcurrentQueue

       public class ProducerConsumerQueue<T> : BlockingCollection<T>
      {
          /// <summary>
          /// Initializes a new instance of the ProducerConsumerQueue, Use Add and TryAdd for Enqueue and TryEnqueue and Take and TryTake for Dequeue and TryDequeue functionality
          /// </summary>
          public ProducerConsumerQueue()  
              : base(new ConcurrentQueue<T>())
          {
          }
      
        /// <summary>
        /// Initializes a new instance of the ProducerConsumerQueue, Use Add and TryAdd for Enqueue and TryEnqueue and Take and TryTake for Dequeue and TryDequeue functionality
        /// </summary>
        /// <param name="maxSize"></param>
          public ProducerConsumerQueue(int maxSize)
              : base(new ConcurrentQueue<T>(), maxSize)
          {
          }
      
      
      
      }
      

      【讨论】:

      • BlockingCollection 默认为队列。所以,我认为这是没有必要的。
      • BlockingCollection 是否像队列一样保留排序?
      • 是的,当它使用 ConcurrentQueue 初始化时
      【解决方案10】:

      从 .NET 5.0/Core 3.0 开始,您可以使用 System.Threading.Channels
      来自this (Asynchronous Producer Consumer Pattern in .NET (C#)) 文章的基准测试显示,与 BlockingCollection 相比,速度显着提升!

      【讨论】:

        猜你喜欢
        • 1970-01-01
        • 2014-10-16
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多