ConcurrentQueue<T>:分段无锁 FIFO 的槽位序号协议
系列:C# 与常用数据结构源码剖析 · 并发集合篇
阅读时间:约 60 分钟
前置知识:FIFO、CAS、内存可见性、线程调度
版本边界:本文讲解当代 .NETConcurrentQueue<T>的公共契约,实现观察固定为dotnet/runtime的v8.0.0/ commit5535e31a712343a63f5d7d796cd874e563e5ac14;主要源码为固定版本的 ConcurrentQueue.cs 与 ConcurrentQueueSegment.cs。私有字段、初始容量、增长上限和快路径都不是公共 API 承诺,其他 runtime tag 必须重新核对。
一、它保证什么,又不保证什么
ConcurrentQueue<T>是支持多生产者、多消费者并发访问的线程安全 FIFO 集合。生产者调用Enqueue,消费者调用TryDequeue;空队列是正常状态,因此取元素用bool而不是抛出异常。
var queue = new ConcurrentQueue<WorkItem>(); queue.Enqueue(item); if (queue.TryDequeue(out WorkItem? work)) Process(work);FIFO 需要精确到并发语义。若入队 A 在入队 B 开始前已完成,后续观测不能把 B 放在 A 前。若两次入队在时间上重叠,则程序没有一个由调用起止时间唯一决定的先后顺序;实现可将它们线性化为任一与真实时间不矛盾的顺序。
线程安全不等于业务操作自动变成事务。例如“如果队列里没有 ID 42 就入队”包含全队列搜索与条件写入,ConcurrentQueue不会把这两步组合成原子操作。同样,它没有容量上限、没有异步等待数据或空位的契约,也不会自动传播完成和取消。
二、为什么“CAS 递增 tail,然后写数组”是错的
一个常见的伪实现是:生产者先用 CAS 将全局tail从i改为i+1,获得槽位i,再写array[i] = item;消费者看到head < tail就读取array[head]。问题在于,生产者的线程可以在 tail 前进后、写 item 前被暂停。消费者会误以为槽位已就绪,读到default(T)或上一轮遗留的值。
Producer P Consumer C CAS tail: i -> i+1 <P 被暂停> 看到 head < tail 读 array[i] // item 尚未发布 array[i] = item只有头尾计数器不足以表达每个槽位处于“空闲、已预留但正在写、可读、已预留但正在读、可再利用”中的哪个阶段。当代 .NET 实现为每个槽位增加序号,将“抢到位置”和“内容已发布”分成两个状态。
三、整体布局:有界环形段串成无界队列
队列由一条 Segment 链组成。每个段是容量固定、且容量适合用位掩码取模的有界多生产者/多消费者环形队列。整体保留头段和尾段引用;尾段写满时冻结它并连接新段,头段排空后前进到下一段。
_head _tail | | v v +------------+ +------------+ +--------------------+ | old segment| ---> | segment | ---> | writable segment | | dequeuing | | full/frozen| | slots + head/tail | +------------+ +------------+ +--------------------+这个设计避免了对一个巨大共享数组做全量扩容复制。段内的高频入队/出队路径使用原子计数器和槽位序号;跨段增长、快照取样等低频操作可以使用一个跨段同步机制。因此“无锁队列”不应被解读为“实现中永远不出现 lock 或等待”;核心单元素快路径的进展性质和所有辅助 API 的实现手段是两个问题。
在v8.0.0的上述两个源文件中,可以看到初始段长度和最大段长度等私有常量,以及在常规增长中扩大后续段的策略。这些数字只可用来解释v8.0.0,不能当作应用程序契约。特别是快照/观察可能使后续增长重新从较小段开始,不能把段容量理解为永不改变的固定 32。
四、槽位序号是一个小型状态机
每个槽位至少包含Item和SequenceNumber。段的逻辑 head/tail 是持续增长的票号,通过掩码映射到环形数组索引。序号不只表示“有/无元素”,还包含槽位属于环形数组哪一轮的代际。
对逻辑入队票号tail映射到的槽位,生产者期待序号表示“该轮可写”。它用 CAS 将 tail 前进,从而独占该逻辑位置;然后写 Item,最后用具有发布语义的序号写入宣布“可读”。消费者只有在序号表示对应 head 已可读时,才能争抢 head 票号并读取 Item。
生产者视角:可写 -> 预留中 -> 写 Item -> 发布为可读 消费者视角:可读 -> 预留中 -> 读 Item -> 清理 -> 发布为下一轮可写下面是结构化伪代码,不是可编译实现,也不是任何 tag 的逐字源码:
// 伪代码:说明发布顺序,省略冻结、溢出和所有边界。 bool TryEnqueue(T item) { while (true) { int tail = Volatile.Read(ref _tail); ref Slot slot = ref _slots[tail & _mask]; int sequence = Volatile.Read(ref slot.Sequence); if (sequence == ExpectedWritableSequence(tail)) { if (Interlocked.CompareExchange(ref _tail, tail + 1, tail) != tail) continue; slot.Item = item; Volatile.Write(ref slot.Sequence, PublishedSequence(tail)); return true; } if (SequenceShowsSegmentFull(sequence, tail)) return false; SpinOrRetry(); } }关键不在函数名,而在状态转移的唯一性:CAS 赢家是唯一可写该逻辑位置的生产者;Item 必须在发布可读序号之前写完;消费者必须在观察到对应序号后才读 Item。这条顺序不能仅依赖“我的 CPU 通常不重排”。
五、Interlocked、Volatile 与 happens-before
Interlocked.CompareExchange既是原子的读-比较-写,又建立必要的内存排序约束。它用于使多个生产者中只有一个取得某个 tail 票号,多个消费者中只有一个取得某个 head 票号。Volatile.Read/Volatile.Write则用于发布和观察状态,避免编译器、JIT 或 CPU 将关键内存操作移到协议不允许的一侧。
在入队路径上,可用下列概念链理解可见性:
Producer: write slot.Item -> release/publish slot.Sequence Consumer: acquire/observe slot.Sequence -> read slot.Item当消费者以对应语义观察到已发布序号时,生产者在发布前的 Item 写入必须对它可见。把序号改成普通读写、把 Item 写放到发布之后,或一厢情愿用Thread.MemoryBarrier散落修补,都很容易破坏证明链。
volatile也不会把x++变成原子操作。两个线程仍可能读到同一旧值,分别写回同一新值。因此票号所有权使用 CAS,状态发布使用 volatile 语义;两者职责不同。
六、入队竞争与线性化点
多个生产者可以同时读到相同 tail,但只有一个 CAS 胜出。失败者重读 tail 和对应槽位序号,它不能在旧槽位上继续写。获得票号的线程即使被暂停,也没有其他生产者会写它的槽位。
但“获得 tail 票号”不能简单当作整个入队已对消费者生效,因为 Item 尚可能没有写完。从可观测行为上,槽位序号的发布使该元素对消费者可见。严格证明线性化点时要结合同一段内前置票号的发布和算法对“空洞”的处理,不要只指着一条 CAS 就宣告证明完成。
当当前段真的无法接受更多入队时,慢路径会重新确认尾段状态,冻结它的入队边界,建立后继段并推进全局尾引用。同步机制必须保证不会由两个线程各自接上一条丢失分支。
七、出队竞争:先确认可读,再争抢 head
消费者不能仅因 tail 看起来领先 head 就读取 Item。它要检查槽位的序号是否已达到当前 head 票号对应的“可读”状态。若某个更早的生产者已预留槽位但尚未发布,后面的生产者即使已写完,消费者也不能跳过前者,否则就破坏 FIFO。
下面同样是教学伪代码:
// 伪代码:不可作为无锁集合实现直接使用。 bool TryDequeue(out T item) { while (true) { int head = Volatile.Read(ref _head); ref Slot slot = ref _slots[head & _mask]; int sequence = Volatile.Read(ref slot.Sequence); if (sequence == ExpectedReadableSequence(head)) { if (Interlocked.CompareExchange(ref _head, head + 1, head) != head) continue; item = slot.Item; if (!PreservedForObservation) { slot.Item = default; Volatile.Write(ref slot.Sequence, NextWritableSequence(head)); } return true; } if (SequenceAndTailShowEmpty(sequence, head)) { item = default; return false; } SpinOrRetry(); } }读取成功后,普通消费路径需要在合适时机清理 Item 引用,并将序号推进到下一轮可写状态。先发布可写、后读或清理 Item 会允许生产者提前覆盖槽位,顺序不可颠倒。
一个失败的TryDequeue只说明该操作线性化时没有可取元素。另一个生产者可以立即入队,因此调用者不能由一次false推导“所有生产者都已完成”。完成协议必须另外设计,例如通道完成、生产者计数或取消令牌。
八、段冻结、链接与容量增长
尾段从可写变为历史段不能只依靠“当前 tail 大于容量”的普通判断。它必须建立一个稳定的入队边界,使新生产者不再进入该段,已经取得票号的生产者仍能完成发布。当代实现使用冻结标记与经调整的 tail 状态达成这个目标;具体偏移量是需按 tag 核验的细节。
跨段增长通常是:发现尾段不能入队,进入慢路径,在跨段同步下二次检查,冻结当前尾段,创建后继段,发布 next 链接,再更新全局 tail 引用。二次检查很重要,因为等待同步的时间里,另一个线程可能已经完成增长。
头段排空且已有后继段时,全局 head 可以前进。旧段何时真正被 GC 回收还取决于是否存在枚举快照等观察者持有它。“从队列头链移除”不等于“立即释放对象”。
九、为什么枚举和 TryPeek 会影响引用清理
ConcurrentQueue<T>支持在并发修改期间获取快照式枚举。为了让枚举器在元素已被其他线程出队后仍能读到快照范围内的 Item,实现不能像普通出队那样立即把所有相关槽位清为default。
在v8.0.0实现中,取快照时会在跨段同步下保留要观察的段,并冻结尾段以获得稳定终点。被标记为观察保留的段,出队路径会保留 Item,直到该段本身不再可达。TryPeek为了在与出队竞争时安全返回实际元素,也可以使当前段进入保留观察状态。
这不是托管内存泄漏:当旧段和所有快照枚举器都不再引用它时,GC 可回收整个段及其保留内容。但它可以表现为短期或中期对象保活,所以在存放大对象、频繁TryPeek、ToArray或长时间持有枚举器的工作负载下,应用内存分析器检查实际保留路径。不要把“TryDequeue 会清理引用”当成每次操作立即清空槽位的绝对承诺。
十、Count、IsEmpty 和 TryPeek 都是观察,不是预留
IsEmpty只回答它进行观察时是否找到元素。下面的检查-操作存在竞争:
if (!queue.IsEmpty) { // 在 IsEmpty 返回后,其他消费者可能已取走元素。 if (queue.TryDequeue(out WorkItem? item)) Process(item); }如果真正意图是“能取就取”,直接调用TryDequeue就是原子边界。TryPeek成功返回的元素也没有被为调用者预留;紧接着的TryDequeue可能由于其他消费者竞争而得到另一结果。因此不能使用 peek 后再 dequeue 实现“只有队首满足条件时才原子移除”。
Count对分段并发结构不是一个简单字段读取。实现需要对头尾段及中间段建立一致的计数视图,并处理冻结与并发前进。具体快路径与同步方式要按 tag 阅读。更重要的是,返回值在调用后就可能过时,不能用Count > 0替代TryDequeue,也不应在每次消费时用 Count 做循环上界。
// 易错:Count 不是对未来出队次数的预留。 for (int i = 0; i < queue.Count; i++) queue.TryDequeue(out _); // 意图清晰:消费直到当前观测为空。 while (queue.TryDequeue(out WorkItem? item)) Process(item);第二段也不意味着队列在循环后永久为空;若生产者仍在运行,它可以马上再添加元素。
十一、无锁不等于无等待、公平或总是更快
lock-free 通常表示在系统整体层面,竞争线程中总有某个操作能在有限步数内取得进展;它不保证每个特定线程都不会饥饿,这是更强的 wait-free 性质所关心的问题。CAS 失败、槽位尚未发布和段切换都可能使线程自旋或重试。
在低竞争、工作量小的场景,Queue<T>加一把锁可能更易理解,并且能把“检查队首条件后移除”等多步业务操作包在同一临界区。无锁结构要支付原子读改写、内存屏障、重试与更复杂的快照协议。不能从“没有全局锁”直接推导延迟或吞吐更优。
严格 FIFO 本身也会建立一个顺序瓶颈:早期生产者抢到槽位后被长时暂停,后续已发布元素不能越过它先出队。这是顺序契约的代价,不是靠多加几次 CAS 就能消除的问题。
十二、ABA 与伪共享:要具体分析,不要贴标签
ABA 指共享位置从 A 变成 B 又回到 A,某线程仅比较当前值时误以为中间没有变化。它在内存手工回收的无锁链表中尤其危险,因为节点地址可被重用。托管引用在仍被线程引用时不会被 GC 回收成另一个对象,而槽位序号又包含环形代际,这些机制具体地处理了“同一索引已绕圈重用”。不应含糊地说“用 CAS 就一定有 ABA”,也不应说“有 GC 就绝对没有 ABA”;必须识别 CAS 比较的具体值与重用协议。
伪共享则是两个逻辑独立的可写字段落在同一缓存行,不同核频繁更新它们时导致缓存一致性流量。head 主要由消费者写,tail 主要由生产者写,实现可通过填充结构使它们分离。但实际效果取决于运行时布局、CPU 缓存行和访问模式,填充还会增加内存。不要根据字段在 C# 声明中相邻就断言已发生伪共享;应阅读固定版本布局并用硬件计数器或基准验证。
十三、与 Channel 和 Queue + lock 的选型
| 需求 | ConcurrentQueue<T> | Channel<T> | Queue<T>+ lock/condition |
|---|---|---|---|
| 立即尝试入/出队 | 适合 | 适合,通过 reader/writer API | 可自行实现 |
| 等待数据到来 | 需额外信号 | 原生支持异步等待 | 需Monitor/信号机制 |
| 有界容量/反压 | 不提供 | 可配置有界通道与满载策略 | 可完全自定义 |
| 完成和异常传播 | 需另外协议 | reader/writer 有完成契约 | 需自己定义 |
| 多步原子业务规则 | 不适合直接组合 | 依 API 契约 | 可在同一锁内实现 |
| 同步、简单的手动轮询 | 适合 | 可用,但能力更完整 | 小型系统容易审查 |
典型误用是用ConcurrentQueue加Task.Delay(1)或忙等实现消费者。延迟轮询在响应时间和空转开销之间两难,取消和完成还要重新造轮子。当需求是异步生产者-消费者流程时,首先评估Channel<T>;当需求是帧开始由主线程快速排空后台结果时,ConcurrentQueue才可能更直接。
Queue<T>加锁不是低级方案。它可以把“查看队首时间、未到期就等待、到期就移除”组合在一个清晰的同步协议中。选型应根据竞争度、异步需求、背压、完成语义和可维护性,而不是将“无锁”当成性能标签。
十四、Unity 中的 Mono、IL2CPP 与平台边界
Unity 项目不能因为代码在当前桌面 .NET 上通过就假设 Editor Mono、Mono Player、IL2CPP Player 和所有 CPU 平台具有完全相同的实现与性能。首先检查目标 Unity Editor 版本的 API Compatibility Level 是否公开所需 API,再对每个目标构建后端和设备做 Player 测试。
Interlocked和 volatile 语义是正确性所需的运行时/平台契约,不能为了“优化 IL2CPP”用普通字段读写替代。但具体指令、屏障成本、自旋策略和托管 GC 保留行为会因 Unity 版本、后端、CPU 和构建配置而异。不要将 CoreCLR 的私有ConcurrentQueue源码逐行套到 Unity 的类库实现上。
在 Unity 中还要区分“线程安全传递数据”与“可以从工作线程调用 Unity API”。将一个结果对象安全入队,不会让该对象内部引用的GameObject、Transform或其他主线程 API 变成可在后台使用。常见模式是后台线程只生成不可变纯数据,主线程在明确的 PlayerLoop 阶段出队并应用到 Unity 对象。
// 后台线程生成纯数据结果。 public readonly record struct MeshChunkResult( int ChunkId, float[] Vertices); // 主线程 Update 阶段调用。 void DrainResults(ConcurrentQueue<MeshChunkResult> results, int budget) { for (int i = 0; i < budget && results.TryDequeue(out var result); i++) ApplyToUnityMesh(result); }上例的budget是每帧工作预算,不是队列容量上限。如果生产速率长期超过消费速率,队列会继续增长;需要有界通道、合并重复任务、丢弃过期结果或反压,而不是只限制每帧出队次数。
十五、故障反例:正确的容器也救不了错误协议
反例一:入队后继续修改对象。队列可以安全发布对象引用,但不会将对象变成不可变快照。生产者入队List<int>后继续添加,消费者同时遍历,仍然是对列表的数据竞争。入队前完成对象构建,后续不再修改,或转移明确所有权。
反例二:用IsEmpty宣告工作完成。队列可以在生产者两次入队之间短暂为空。需要单独的“所有生产者已结束”信号,并在该信号成立且当前排空后结束消费。
反例三:忙等直到有数据。不受控的while (!TryDequeue) {}会持续消耗 CPU,单核环境下还可能影响生产者获得调度。需要等待时使用 Channel 或与队列正确组合的信号机制,但信号计数和入队必须遵守一致顺序,避免丢失唤醒。
反例四:无限生产者遇到慢消费者。ConcurrentQueue能安全增长,不代表内存无限。应监控排队延迟、生产/消费速率与内存,并在系统设计中定义过载策略。
反例五:遍历队列寻找任务后假设它仍在队中。枚举是快照观察,不是锁定或句柄。若需要按 ID 取消、更新或去重,要建立额外索引/状态机,或在单所有者线程内串行化所有变更。
十六、可复现压力测试
无锁错误往往不会在一次运行中出现。测试应分为契约测试、受控交错和长时间压力三层。对 BCL 队列的应用验证,重点是自己的所有权和完成协议;对自研队列,还要验证每一个内存序和槽位代际不变式。
一个基础 MPMC 压力测试可以让P个生产者各自生成不重复的(producerId, sequence),C个消费者并发出队。结束时验证:总数恰好相等;每个 ID 恰好出现一次;对同一生产者,若它的入队调用按程序顺序串行完成,出队观测序号不能倒退。不要对不同生产者的重叠调用强加未承诺的全局时间顺序。
public readonly record struct Message(int Producer, int Sequence); // 结构化测试框架:省略取消、异常聚合和超时处理。 var queue = new ConcurrentQueue<Message>(); var consumed = new ConcurrentBag<Message>(); int producersRemaining = producerCount; Task[] producers = Enumerable.Range(0, producerCount).Select(id => Task.Run(() => { for (int sequence = 0; sequence < itemsPerProducer; sequence++) queue.Enqueue(new Message(id, sequence)); Interlocked.Decrement(ref producersRemaining); })).ToArray(); Task[] consumers = Enumerable.Range(0, consumerCount).Select(_ => Task.Run(() => { while (Volatile.Read(ref producersRemaining) != 0 || !queue.IsEmpty) { if (queue.TryDequeue(out Message message)) consumed.Add(message); else Thread.Yield(); } })).ToArray(); await Task.WhenAll(producers.Concat(consumers));这是测试骨架,不是生产级完成协议。它需要超时防止 CI 永久挂起,需要捕获所有任务异常,并在测试后按 ID 验证丢失、重复和生产者内部顺序。为了增加交错多样性,可使用固定随机种子在操作间插入Thread.Yield、小量自旋或可控屏障,并将失败种子与完整环境保存下来。
还应专门覆盖:频繁跨段;生产者获得票号后暂停;队列在空与非空之间高频切换;与TryPeek、Count、ToArray和枚举并发;元素为引用类型且检查最终可回收性;计数器接近边界时的代际运算。最后一项通常需要对自研实现提供小位宽测试模型,而不是在真实int计数器上等待极长时间。
基准测试和正确性压力测试应分开。基准需固定 runtime 版本、CPU/核数、生产者与消费者数、元素大小、队列稳态深度和内存限制,分别测 1P1C、MP1C、1PMC 和 MPMC。要同时记录吞吐、延迟分布、CPU 使用和分配,不能只发布一个没有原始报告的“快几倍”。
十七、源码审查清单
阅读一个具体版本时,可按以下问题逐层核对:
- 当前文件属于哪个 runtime tag/commit,与本地实际运行的版本是否匹配?
- Segment 容量、掩码、head/tail 填充布局与槽位序号初值如何定义?
- 生产者何时独占票号,何时写 Item,何时发布可读?
- 消费者如何区分“真空”和“更早生产者尚未发布”?
- 槽位完成出队后,Item 与下一轮可写序号按什么顺序更新?
- 段冻结如何阻止新入队,又不破坏已预留槽位的发布?
- 新段容量如何选择,为什么观察保留会影响后续增长?
TryPeek、枚举、ToArray和Count如何获得可观测的视图,哪些路径使用跨段同步?- 被观察段为什么不清理 Item,它最终如何变得不可达?
- 哪些语句是原子所有权竞争,哪些是状态发布,它们共同建立了哪条 happens-before 链?
结语
ConcurrentQueue<T>不是“把 Queue 的 head 和 tail 换成 CAS”。当代实现将多个有界环形段串起来,每个槽位用带代际的序号区分预留和发布,CAS 决定票号所有权,volatile 语义建立 Item 的可见性。段冻结和链接解决增长,观察保留协议则在并发出队下支撑快照枚举。
使用时要记住:IsEmpty、Count和TryPeek都只是观察,不为下一步预留状态;队列线程安全不意味元素对象可被并发修改;无界队列没有反压;无锁也不意味公平、无自旋或必然更快。需要异步等待、完成和有界容量时优先评估Channel<T>;需要多步原子业务规则时,Queue<T>加明确锁协议可能更正确。
最后,这类源码只能在固定 runtime tag 下逐行讨论。将公共 FIFO 契约、特定版本私有布局和用于解释的伪代码分开,才能真正理解算法,而不是将一个简化示例误当成可投产的无锁队列。
下一篇:ConcurrentStack<T>:无锁链式栈与 CAS