【问题标题】:BlockingCollection Max SizeBlockingCollection 最大尺寸
【发布时间】:2013-04-10 14:26:48
【问题描述】:

我了解使用 ConcurrentQueue 的 BlockingCollection 的有限容量为 100。

但是我不确定这意味着什么。

我正在尝试实现一个并发缓存,如果队列大小太大,它可以出队,可以在一个操作中出队/入队(即缓存溢出时松散的消息)。有没有办法为此使用有限容量,或者手动执行此操作或创建一个新集合更好。

基本上我有一个阅读线程和几个写作线程。如果队列中的数据是所有作者中“最新鲜”的,我希望它。

【问题讨论】:

  • 您可能需要某种后台事件来告知作者最新数据?

标签: c# .net-4.0 concurrency


【解决方案1】:

N 的有限容量意味着如果队列已经包含 N 个项目,则任何尝试添加另一个项目的线程都会阻塞,直到另一个线程移除一个项目。

您似乎想要的是一个不同的概念 - 您希望最近添加的项目成为消费线程出列的第一个项目。

您可以通过使用ConcurrentStack 而不是底层存储的 ConcurrentQueue 来实现。

您可以使用this constructor 并传入ConcurrentStack

例如:

var blockingCollection = new BlockingCollection<int>(new ConcurrentStack<int>());

通过使用ConcurrentStack,您可以确保消费线程出队的每个项目都是当时队列中最新鲜的项目。

另请注意,如果您为阻塞集合指定上限,您可以使用BlockingCollection.TryAdd(),如果您调用它时集合已满,则返回false

【讨论】:

  • 基于我最新的真实陈述。但是我仍然希望队列按照它们入队的顺序进行调度。似乎最好在每次添加之前执行 if(q.count>x) {q.take();//donothing} q.add(newitem) 以减小大小,但是如果另一个线程在之后长度被计算然后你最终不必要地删除了一条消息
  • 嗯我不认为我真的理解你的要求......但是如果你为阻塞集合指定一个上限,你可以使用BlockingCollection.TryAdd(),如果队列是满 - 你能用吗?
【解决方案2】:

在我看来,您正在尝试构建 MRU(最近使用的)缓存之类的东西。 BlockingCollection 不是最好的方法。

我建议您改用LinkedList。它不是线程安全的,因此您必须提供自己的同步,但这并不太难。您的入队方法如下所示:

LinkedList<MyType> TheQueue = new LinkedList<MyType>();
object listLock = new object();

void Enqueue(MyType item)
{
    lock (listLock)
    {
        TheQueue.AddFirst(item);
        while (TheQueue.Count > MaxQueueSize)
        {
            // Queue overflow. Reduce to max size.
            TheQueue.RemoveLast();
        }
    }
}

而且出队更容易:

MyType Dequeue()
{
    lock (listLock)
    {
        return (TheQueue.Count > 0) ? TheQueue.RemoveLast() : null;
    }
}

如果您希望消费者在队列中进行非忙碌等待,则涉及更多一些。您可以使用Monitor.WaitMonitor.Pulse 来实现。有关示例,请参见 Monitor.Pulse 页面上的示例。

更新:

我突然想到,你可以用循环缓冲区(数组)做同样的事情。只需保持头部和尾部指针。您在head 插入并在tail 删除。如果你去插入,和head == tail,那么你需要增加tail,这有效地删除了之前的tail项。

【讨论】:

    【解决方案3】:

    如果您想要一个自定义的BlockingCollection,它包含 N 个最近的元素,并在其已满时删除最旧的元素,您可以基于 Channel&lt;T&gt; 轻松创建一个。 Channels 旨在用于异步场景,但让它们阻塞消费者是微不足道的,并且不应导致任何不需要的副作用(如死锁),即使在安装了 SynchronizationContext 的环境中使用也是如此。

    public class MostRecentBlockingCollection<T>
    {
        private readonly Channel<T> _channel;
    
        public MostRecentBlockingCollection(int capacity)
        {
            _channel = Channel.CreateBounded<T>(new BoundedChannelOptions(capacity)
            {
                FullMode = BoundedChannelFullMode.DropOldest,
            });
        }
    
        public bool IsCompleted => _channel.Reader.Completion.IsCompleted;
    
        public void Add(T item)
            => _channel.Writer.WriteAsync(item).AsTask().GetAwaiter().GetResult();
    
        public T Take()
            => _channel.Reader.ReadAsync().AsTask().GetAwaiter().GetResult();
    
        public void CompleteAdding() => _channel.Writer.Complete();
    
        public IEnumerable<T> GetConsumingEnumerable()
        {
            while (_channel.Reader.WaitToReadAsync().AsTask().GetAwaiter().GetResult())
                while (_channel.Reader.TryRead(out var item))
                    yield return item;
        }
    }
    

    MostRecentBlockingCollection 类只阻止消费者。生产者总是可以在集合中添加项目,从而(可能)导致一些以前添加的元素被删除。

    添加取消支持应该很简单,因为Channel&lt;T&gt; API 已经支持它。添加对超时的支持不是那么简单,但应该不会很难做到。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2014-04-26
      • 1970-01-01
      • 1970-01-01
      • 1970-01-01
      • 2018-04-21
      • 2015-02-10
      相关资源
      最近更新 更多