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

多生产者入队、单消费者批量出队的ConcurrentQueue同步方案选型

解决方案:使用lock配合普通Queue<T>

你的核心需求是在Flush操作期间完全阻止生产者入队,同时一次性取出所有队列元素后释放锁,避免耗时的DoSomething阻塞生产者。结合这个场景,最直接且高效的方案是放弃ConcurrentQueue<T>,改用普通Queue<T>配合lock同步,原因如下:

  • ConcurrentQueue<T>的线程安全特性依赖内部CAS操作,但你需要的是独占式访问(Flush时禁止所有入队),此时ConcurrentQueue的内部并发控制完全多余,反而增加不必要的性能开销。
  • lock语法简单直观,能完美实现"生产者互斥入队"、"Flush独占队列"的需求,且在单个消费者、多生产者的场景下性能足够。

最终实现代码

public class Cache
{
    // 使用普通Queue替代ConcurrentQueue
    private readonly Queue<Item> _queue = new();
    // 专用锁对象,避免使用this或其他公开对象引发意外锁竞争
    private readonly object _queueLock = new();

    public void Enqueue(Item item)
    {
        ArgumentNullException.ThrowIfNull(item, nameof(item));
        lock (_queueLock)
        {
            _queue.Enqueue(item);
        }
    }

    public void Flush()
    {
        List<Item> itemsToProcess;
        // 仅在获取队列元素的阶段持有锁
        lock (_queueLock)
        {
            // 一次性复制所有元素并清空队列
            itemsToProcess = _queue.ToList();
            _queue.Clear();
        }

        // 释放锁后再执行耗时的处理逻辑,不阻塞生产者入队
        foreach (var item in itemsToProcess)
        {
            DoSomething(item);
        }
    }

    public void DoSomething(Item item)
    {
        // ... 你的业务逻辑
    }
}

方案优势

  1. 严格满足需求:
    • 多个生产者调用Enqueue时,lock保证同一时间只有一个生产者能修改队列;
    • Flush执行时,lock会阻止所有生产者入队,直到队列元素被全部取出并清空,锁才会释放。
  2. 性能优化:
    • 耗时的DoSomething在锁外执行,生产者无需等待处理完成即可继续入队,最大化并发效率;
    • 普通Queue<T>的操作比ConcurrentQueue<T>更轻量,配合lock的开销远低于ConcurrentQueue的内部并发控制。

其他可选方案(不推荐)

如果坚持使用ConcurrentQueue<T>,也可以给Enqueue和Flush的队列操作加lock,但这样会导致ConcurrentQueue的内部CAS操作与外部lock双重开销,完全没有必要,代码示例仅供参考:

public class Cache
{
    private readonly ConcurrentQueue<Item> _queue = new();
    private readonly object _queueLock = new();

    public void Enqueue(Item item)
    {
        ArgumentNullException.ThrowIfNull(item, nameof(item));
        lock (_queueLock)
        {
            _queue.Enqueue(item);
        }
    }

    public void Flush()
    {
        var itemsToProcess = new List<Item>();
        lock (_queueLock)
        {
            while (_queue.TryDequeue(out var item))
            {
                itemsToProcess.Add(item);
            }
        }

        foreach (var item in itemsToProcess)
        {
            DoSomething(item);
        }
    }

    public void DoSomething(Item item)
    {
        // ...
    }
}

内容的提问来源于stack exchange,提问作者lonix

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.15 21:42:51