src/DotNetCommon.Core/ProducerConsumerQueue.cs

(开头部分) 12KB

这里只显示每个文件的开头 60 行。登录后可以解锁完整代码。

namespace DotNetCommon
{
    using System;
    using System.Collections.Concurrent;
    using System.Runtime.Serialization;
    using System.Threading;
    using System.Threading.Tasks;

    /// <summary>
    /// 使用并行编程技术 (<c>TPL</c>) 实现的一个消息队列模型(<c>生产者/消费者</c>), 基于 <c>System.Collections.Concurrent.BlockingCollection&lt;T></c>
    /// <example>
    /// 下面是封装的非阻塞日志: 
    /// <code>
    /// public class Logger<para/>
    /// {<para/>
    ///    private static readonly ProducerConsumerQueue&lt;string> logQueue = new ProducerConsumerQueue&lt;string>(msg =><para/>
    ///     {<para/>
    ///         File.AppendAllText("d:/2020-10-10.log", msg);<para/>
    ///     }, 1, 10000);<para/>
    ///     public static void Log(string msg)<para/>
    ///     {<para/>
    ///         logQueue.Add(msg);<para/>
    ///     }<para/>
    /// }<para/>
    /// </code>
    /// </example>
    /// </summary>
    /// <typeparam name="T">消息队列中存储的模型</typeparam>
    public sealed class ProducerConsumerQueue<T> : IDisposable
    {
        private readonly BlockingCollection<T> _queue;

        /// <summary>
        /// 根据指定的消费者逻辑、最大并发数,创建一个不限消息数量的队列模型
        /// </summary>
        /// <param name="consumer">消费者逻辑</param>
        /// <param name="maxConcurrencyLevel">最多有多少个消费者</param>
        public ProducerConsumerQueue(Action<T> consumer, uint maxConcurrencyLevel)
            : this(consumer, maxConcurrencyLevel, -1) { }

        /// <summary>
        /// 根据指定的消费者逻辑、最大并发数以及最大的消息数量,创建一个队列模型
        /// </summary>
        /// <param name="consumer">消费者逻辑</param>
        /// <param name="maxConcurrencyLevel">最多有多少个消费者</param>
        /// <param name="boundedCapacity">最多存储多少条消息</param>
        public ProducerConsumerQueue(Action<T> consumer, uint maxConcurrencyLevel, uint boundedCapacity)
            : this(consumer, maxConcurrencyLevel, (int)boundedCapacity) { }

        private ProducerConsumerQueue(Action<T> consumer, uint maxConcurrencyLevel, int boundedCapacity)
        {
            Ensure.NotNull(consumer, nameof(consumer));
            Ensure.That(maxConcurrencyLevel > 0, $"{nameof(maxConcurrencyLevel)} should be greater than zero.");
            Ensure.That(boundedCapacity != 0, $"{nameof(boundedCapacity)} should be greater than zero.");

            _queue = boundedCapacity < 0 ? new BlockingCollection<T>() : new BlockingCollection<T>(boundedCapacity);

            MaximumConcurrencyLevel = maxConcurrencyLevel;
            Completion = Configure(consumer);
        }
后面还有 230 行代码,解锁后查看完整代码

24 小时内免费解锁 3 个项目,之后 1 积分/个。 规则说明

AI 解读

登录后可用,每次 10 积分,解读结果公开显示在下面。

还没有人解读过这个文件。