You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

.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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.07.20 10:09:57