.NET Framework 4.8下带优先级的生产者消费者模式最优实现问询
带优先级的生产者消费者队列实现建议(.NET Framework 4.8)
需求概述
- 多线程可安全入队元素
- 单线程消费元素
- 维持同优先级内的入队顺序
- 优先级数量有限
现有实现及问题
BlockingCollection实现
public class PriorityQueue : IProducerConsumer { private readonly BlockingCollection<Action> _lowQueue; private readonly BlockingCollection<Action> _normalQueue; private readonly BlockingCollection<Action> _highQueue; private readonly CancellationToken _cancelToken; private Task _dequeTask; public PriorityQueue(CancellationToken cancelToken) { _lowQueue = new BlockingCollection<Action>(new ConcurrentQueue<Action>()); _normalQueue = new BlockingCollection<Action>(new ConcurrentQueue<Action>()); _highQueue = new BlockingCollection<Action>(new ConcurrentQueue<Action>()); } public void Enqueue(Action item, QueuePriority priority = QueuePriority.Default) { BlockingCollection<Action> queue; switch (priority) { case QueuePriority.Low: queue = _lowQueue; break; case QueuePriority.Normal: queue = _normalQueue; break; case QueuePriority.Hight: queue = _highQueue; break; default: queue = _lowQueue; break; } if (!queue.IsAddingCompleted) { queue.Add(item); } } public void Stop() { _lowQueue.CompleteAdding(); _normalQueue.CompleteAdding(); _highQueue.CompleteAdding(); } public async Task StopAsync() { await Task.Run(() => Stop()); } public void StartDequeing() { _dequeTask = Task.Factory.StartNew(DequeueTask, _cancelToken, TaskCreationOptions.LongRunning, TaskScheduler.Default); } private void DequeueTask() { var queues = new BlockingCollection<Action>[] { _highQueue, _normalQueue, _lowQueue }; while (!_cancelToken.IsCancellationRequested) { BlockingCollection<Action>.TakeFromAny(queues, out Action item, _cancelToken); item?.Invoke(); } } }
TPL DataFlow实现(已排除)
该方案通过BufferBlock转发元素实现优先级,但经基准测试性能表现较差,因此不再考虑。
public class PriorityFlowQueue : IProducerConsumer { private readonly BufferBlock<Action> _lowBuffer; private readonly BufferBlock<Action> _normalBuffer; private readonly BufferBlock<Action> _highBuffer; private readonly BufferBlock<Action> _sourceBuffer; private readonly CancellationToken _cancelToken; private Task _dequeTask; public PriorityFlowQueue(CancellationToken cancelToken) { _cancelToken = cancelToken; var options = new DataflowBlockOptions() { EnsureOrdered = true }; _lowBuffer = new BufferBlock<Action>(options); _normalBuffer = new BufferBlock<Action>(options); _highBuffer = new BufferBlock<Action>(options); _sourceBuffer = new BufferBlock<Action>(options); Task.Run(ForwardToSource, _cancelToken); } public void Enqueue(Action item, QueuePriority priority = QueuePriority.Default) { var buffer = GetBufferBlockByPriority(priority); buffer.Post(item); } public async Task EnqueueAsync(Action item, QueuePriority priority = QueuePriority.Default) { var buffer = GetBufferBlockByPriority(priority); await buffer.SendAsync(item); } public void StartDequeing() { _dequeTask = Task.Run(DequeueTask, _cancelToken); } public async Task StopAsync() { _highBuffer.Complete(); _normalBuffer.Complete(); _lowBuffer.Complete(); await Task.WhenAll(_highBuffer.Completion, _normalBuffer.Completion, _lowBuffer.Completion); } private async Task ForwardToSource() { while (!_cancelToken.IsCancellationRequested) { await Task.WhenAny(_highBuffer.OutputAvailableAsync(), _normalBuffer.OutputAvailableAsync(), _lowBuffer.OutputAvailableAsync()); Action item; if (_highBuffer.TryReceive(out item)) { } else if (_normalBuffer.TryReceive(out item)) { } else if (_lowBuffer.TryReceive(out item)) { } await _sourceBuffer.SendAsync(item); } } private async Task DequeueTask() { while (await _sourceBuffer.OutputAvailableAsync()) { var item = await _sourceBuffer.ReceiveAsync(); item?.Invoke(); } } private BufferBlock<Action> GetBufferBlockByPriority(QueuePriority priority) { BufferBlock<Action> buffer; switch (priority) { case QueuePriority.Low: buffer = _lowBuffer; break; case QueuePriority.Normal: buffer = _normalBuffer; break; case QueuePriority.Hight: buffer = _highBuffer; break; default: buffer = _lowBuffer; break; } return buffer; } }
优化实现建议
1. 优化现有BlockingCollection实现
你的现有方案已经契合需求,可通过以下调整提升健壮性与性能:
- 修正优先级枚举拼写错误(
Hight改为High),避免逻辑错误 - 复用优先级队列数组,避免循环内重复创建
- 捕获
TakeFromAny抛出的取消异常,防止任务意外终止 - 在停止时等待消费任务完成,避免线程泄漏
优化后的核心代码片段:
private readonly BlockingCollection<Action>[] _priorityQueues; private readonly CancellationToken _cancelToken; private Task _dequeTask; public PriorityQueue(CancellationToken cancelToken) { _cancelToken = cancelToken; var highQueue = new BlockingCollection<Action>(new ConcurrentQueue<Action>()); var normalQueue = new BlockingCollection<Action>(new ConcurrentQueue<Action>()); var lowQueue = new BlockingCollection<Action>(new ConcurrentQueue<Action>()); _priorityQueues = new[] { highQueue, normalQueue, lowQueue }; } private void DequeueTask() { try { while (!_cancelToken.IsCancellationRequested) { int queueIndex = BlockingCollection<Action>.TakeFromAny(_priorityQueues, out Action item, _cancelToken); if (queueIndex >= 0) { item?.Invoke(); } } } catch (OperationCanceledException) { // 预期的取消逻辑,无需额外处理 } finally { foreach (var queue in _priorityQueues) { queue.CompleteAdding(); } } } public async Task StopAsync() { foreach (var queue in _priorityQueues) { queue.CompleteAdding(); } if (_dequeTask != null) { await _dequeTask.ConfigureAwait(false); } }
2. 基于ConcurrentQueue+ManualResetEventSlim的轻量实现
若追求更优性能,可直接使用ConcurrentQueue配合信号量手动实现,减少BlockingCollection的封装开销:
public class CustomPriorityQueue : IProducerConsumer { private readonly ConcurrentQueue<Action> _highQueue = new ConcurrentQueue<Action>(); private readonly ConcurrentQueue<Action> _normalQueue = new ConcurrentQueue<Action>(); private readonly ConcurrentQueue<Action> _lowQueue = new ConcurrentQueue<Action>(); private readonly ManualResetEventSlim _hasItemsSignal = new ManualResetEventSlim(false); private readonly CancellationToken _cancelToken; private Task _consumeTask; private bool _isStopped; public CustomPriorityQueue(CancellationToken cancelToken) { _cancelToken = cancelToken; } public void Enqueue(Action item, QueuePriority priority = QueuePriority.Default) { if (_isStopped) throw new InvalidOperationException("队列已停止接受新元素"); switch (priority) { case QueuePriority.High: _highQueue.Enqueue(item); break; case QueuePriority.Normal: _normalQueue.Enqueue(item); break; case QueuePriority.Low: _lowQueue.Enqueue(item); break; default: _lowQueue.Enqueue(item); break; } _hasItemsSignal.Set(); } public void StartDequeing() { _consumeTask = Task.Factory.StartNew(ConsumeLoop, _cancelToken, TaskCreationOptions.LongRunning, TaskScheduler.Default); } private void ConsumeLoop() { while (!_cancelToken.IsCancellationRequested && !_isStopped) { _hasItemsSignal.Wait(_cancelToken); // 按优先级顺序尝试消费 if (_highQueue.TryDequeue(out var highItem)) { highItem?.Invoke(); } else if (_normalQueue.TryDequeue(out var normalItem)) { normalItem?.Invoke(); } else if (_lowQueue.TryDequeue(out var lowItem)) { lowItem?.Invoke(); } else { _hasItemsSignal.Reset(); } } } public void Stop() { _isStopped = true; _hasItemsSignal.Set(); } public async Task StopAsync() { Stop(); if (_consumeTask != null) { await _consumeTask.ConfigureAwait(false); } _hasItemsSignal.Dispose(); } }
该实现通过手动控制信号量减少锁开销,同时保留了ConcurrentQueue的线程安全性和FIFO顺序特性,性能更优。
关键注意事项
- 同优先级顺序:
ConcurrentQueue天然保证FIFO,满足同优先级元素的入队顺序要求 - 线程安全:
ConcurrentQueue的入队、出队操作均为线程安全,多线程入队无需额外加锁 - 资源清理:停止时需唤醒消费线程、等待任务完成,并释放信号量等资源
内容的提问来源于stack exchange,提问作者JuanDYB
相关产品推荐
相关产品推荐

