【问题标题】:Dequeue from Queue with where expression使用 where 表达式从队列中出列
【发布时间】:2017-11-08 08:24:43
【问题描述】:

我有一个基于 Marc Gravell's Blocking Queue 的阻塞队列类

现在我有一些情况,我只想让某个对象出队,我知道这不是队列的真正用例,但在某些情况下,我认为这是一个很好的扩展,例如等待某个网络答案.

这有点像

TryDequeueWhere(Func<T, bool> expression, out T value, int? waitTimeInMs = null)

问题是我不知道如何等待和阻止某个对象。

【问题讨论】:

  • 我建议使用 BlockingCollectionBlockingQueue
  • 如果你打电话给TryDequeueWhere会发生什么?假设匹配条目位于索引 5。究竟会发生什么?索引 6 处的条目会移动到 5(等等)吗?
  • 如果某个元素与表达式匹配,则在该位置移除该元素,因此下一个元素将向右上移。
  • 简短的回答是这不是队列(如您所知)。你能建吗?当然。基本上你需要锁定整个队列,foreach 直到你找到匹配的元素。如果找到匹配项,则需要新建一个全新的队列(包含匹配项之前和之后的元素)。然后解锁。

标签: c# queue mutex blockingqueue


【解决方案1】:

codereview 上发布我的代码并根据其他用户的建议进行改进后(感谢Pieter Witvoet)这是我的最终代码

using System;
using System.Collections.Generic;
using System.Linq;
using System.Threading;

public class BlockingQueue<T>: IDisposable
{

    /// <summary>   The queue based on a list, to extract from position and remove at position. </summary>
    private readonly List<QueueObject<T>> queue = new List<QueueObject<T>>();
    private bool _closing;


    private class QueueObject<T>
    {
        //// <summary>   Constructor. </summary>
        /// <param name="timeStamp">    The time stamp when the object is enqueued. </param>
        /// <param name="queuedObject"> The queued object. </param>
        public QueueObject(DateTime timeStamp, T queuedObject)
        {
            TimeStamp = timeStamp;
            QueuedObject = queuedObject;
        }

        /// <summary>   Gets or sets the queued object. </summary>
        /// <value> The queued object. </value>
        public T QueuedObject { get; private set; }

        /// <summary>   Gets or sets timestamp, when the object was enqueued. </summary>
        /// <value> The time stamp. </value>
        public DateTime TimeStamp { get; private set; }
    }


    public void Enqueue(T item)
    {
        lock (queue)
        {
            // Add an object with current time to the queue
            queue.Add(new QueueObject<T>(DateTime.Now, item));


            if (queue.Count >= 1)
            {
                // wake up any blocked dequeue
                Monitor.PulseAll(queue);
            }
        }
    }

    /// <summary>   Try dequeue an object that matches the passed expression. </summary>
    /// <param name="expression">   The expression that an object has to match. </param>
    /// <param name="value">        [out] The resulting object. </param>
    /// <param name="waitTimeInMs"> (Optional)  The time in ms to wait for the item to be returned. </param>
    /// <returns>   An object that matches the passed expression. </returns>
    public bool TryDequeueWhere(Func<T, bool> expression, out T value, int? waitTimeInMs = null)
    {
        // Save the current time to later calculate a new timeout, if an object is enqueued and does not match the expression.
        DateTime dequeueTime = DateTime.Now;
        lock (queue)
        {
            while (!_closing)
            {
                if (waitTimeInMs == null)
                {
                    while (queue.Count == 0)
                    {
                        if (_closing)
                        {
                            value = default(T);
                            return false;
                        }
                        Monitor.Wait(queue);
                    }
                }
                else
                {
                    // Releases the lock on queue and blocks the current thread until it reacquires the lock. 
                    // If the specified time-out interval elapses, the thread enters the ready queue.
                    if (!Monitor.Wait(queue, waitTimeInMs.Value))
                    {
                        break;
                    }
                    try
                    {
                        // select the object by the passed expression
                        var queuedObjects = queue.Select(q => q.QueuedObject).ToList();
                        // Convert the expression to a predicate to get the index of the item
                        Predicate<T> pred = expression.Invoke;
                        int indexOfQueuedObject = queuedObjects.FindIndex(pred);
                        // if item is found, get it and remove it from the list
                        if (indexOfQueuedObject >= 0)
                        {
                            value = queuedObjects.FirstOrDefault(expression);
                            queue.RemoveAt(indexOfQueuedObject);
                            return true;
                        }
                    }
                    catch (Exception)
                    {
                        break;
                    }
                    // If item was not found, calculate the remaining time and try again if time is not elapsed.
                    var elapsedTime = (DateTime.Now - dequeueTime).TotalMilliseconds;
                    if ((int) elapsedTime >= waitTimeInMs.Value)
                    {
                        break;
                    }
                    waitTimeInMs = waitTimeInMs.Value - (int) elapsedTime;
                }
            }
        }
        value = default(T);
        return false;
    }

    /// <summary> Close the queue and let finish all waiting threads. </summary>
    public void Close()
    {
        lock (queue)
        {
            _closing = true;
            Monitor.PulseAll(queue);
        }
    }

    /// <summary>
    /// Performs application-defined tasks associated with freeing, releasing, or resetting unmanaged
    /// resources.
    /// </summary>
    public void Dispose()
    {
        Close();
    }

}

【讨论】:

    猜你喜欢
    • 2011-09-26
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2013-04-28
    • 1970-01-01
    • 1970-01-01
    • 2022-12-20
    • 2012-05-22
    相关资源
    最近更新 更多