核心思想

        用原子性序号(sequence / head / tail)和固定数组槽(slot)来把“谁可以写 / 谁可以读”这个互斥语义用数值顺序表达出来,CAS只需要保证只有一个线程可以把序号推进,从而代替传统锁的互斥与同步。

        关键点在于不把整个临界区锁住,而是将单独冲突点现在在一个或者几个插槽中,让这几个goroutine进行竞争这些个插槽,而不影响其他插槽正常使用,这会大大的提升锁的细粒度,提高性能。

核心组成部分
这种无锁的高性能思路在 DisruptorRistrettoGo runtime 的 netpollKCP 协议 等系统里广泛应用。

一般结构
RingBuffer {
    padding
    capacity (power of 2)
    mask = capacity - 1
    slots[]            // fixed array of Slot
    padding
    head (consumer sequence)   // atomic uint64
    padding
    tail (producer sequence)   // atomic uint64
    padding
}
Slot {
    // optionally: sequence / state / flag (atomic)
    value
    padding to avoid false sharing
}

关键设计点:

  • capacity 为 2^n:取模使用 idx = seq & mask,比 % 快。

  • cache-line paddingheadtail 必须分开位于不同缓存行,避免 false sharing。

  • slot 中可带 sequence 或状态位(见下文 MPMC 方案),用于标记该槽是否已准备好供消费。

简单入门:SPSC

        单一生产者和消费者的情况很简单,无需slot插槽的概念,只需要tail和head即可。下面是简单伪代码实现:

type RingBuffer struct {
    buf []T
    mask uint64
    head uint64 // consumer index: next to read
    tail uint64 // producer index: next to write
}

func New(capacity int) *RingBuffer { ... } // capacity power of two

// Producer (single)
func (r *RingBuffer) Enqueue(v T) bool {
    for {
        tail := atomic.LoadUint64(&r.tail)
        head := atomic.LoadUint64(&r.head) // safe: single consumer updates head
        if tail - head == uint64(len(r.buf)) {  //头-尾等于一个环长度
            return false // full
        }
        // slot index:
        idx := tail & r.mask
        r.buf[idx] = v                   // write payload
        // publish by advancing tail atomically (no other producer)
        atomic.StoreUint64(&r.tail, tail+1)
        return true
    }
}

// Consumer (single)
func (r *RingBuffer) Dequeue() (T, bool) {
    head := atomic.LoadUint64(&r.head)
    tail := atomic.LoadUint64(&r.tail)
    if head == tail {
        return nil, false // empty
    }
    idx := head & r.mask
    v := r.buf[idx]
    atomic.StoreUint64(&r.head, head+1)
    return v, true
}

SPSC 不需要 CAS,因为只有一个写序列,一个读序列,atomic.Store+Load 保证可见性。

进阶:MPMC

        多消费者和多生产者往往才是这类问题的核心难点,主流实现方法是:全局Tail+slot状态,核心思路如下:

  • 生产者用 seq = atomic.FetchAdd(&tail, 1) 或者 CAS 尝试获得一个 seq(取号)。

  • 但是要确保不会覆盖尚未被消费的槽(即 seq - head < capacity 校验)。

  • 为了区分槽被消费还是尚存旧数据,每个 slot 维护一个 slot.sequence类似于版本号),用于判定该 slot 当前属于哪个 sequence。

    • 生产者写前检查 slot.sequence == seq(表示可写);写入后把 slot.sequence = seq+1(或其它标记)表示已发布。

    • 消费者读取 slot 时根据 slot.sequence 判断是否已发布。

结构定义(实际上使用的是一个队列类型结构模型)
// 每个槽位结构体,保存该槽的当前状态(sequence)以及数据内容(value)
type Slot struct {
    sequence uint64 // 原子变量:标识该槽位当前属于哪个逻辑序号(决定是否可读/可写)
    value    T      // 存储实际的数据
    // pad 是为了防止 false sharing(不同线程频繁读写相邻变量导致缓存行抖动)
    pad [56]byte    
}

// 整个环形缓冲区结构
type RingBuffer struct {
    slots []Slot   // 固定大小的槽位数组(容量通常是 2 的幂次)
    mask  uint64   // 用于取模优化(idx = seq & mask)
    tail  uint64   // 原子变量:生产者下一个要写入的序号(next sequence to produce)
    head  uint64   // 原子变量:消费者下一个要读取的序号(next sequence to consume)
}

// 初始化函数
func New(capacity int) *RingBuffer {
    // 通常 capacity 要取 2 的幂次方,如 1024, 2048...
    rb := &RingBuffer{
        slots: make([]Slot, capacity),
        mask:  uint64(capacity - 1),
    }
    // 初始化每个槽位的 sequence
    // 理解:每个槽位起始时认为“已经被消费掉”,即 sequence == i(空闲可写)
    for i := 0; i < capacity; i++ {
        rb.slots[i].sequence = uint64(i)
    }
    return rb
}

【重点】

    slots是当前存储结构的下标,但是这个下标存储的实际数据不只有所需要存储的数据,还有一个标记值,这个标记值sequence在每一个写操作后都会将当前值+1处理,这样就是为了标识当前插槽块被修改过,其实就是类似于版本号,这样可以保证高并发情况下,操作在修改插槽内数据时候,会在写入slot中先进行对比当前sequence是否跟操作执行前记录的值是一样的,如果一样则允许修改,不一样则表示当前数据被其他协程改动过,不允许修改。

        这其实就是实现了同CAS一样的思想,也是Java中的CurrentedHashMap一样,或者说List的RandomAccess一样。

【问题】为什么需要CAS操作slot,还需要sequence进行类CAS操作记录状态?

        因为CAS操作仅仅针对同一个slot,但是具体这个slot内数据是否已经被消费过,这个就需要sequence进行辅助判断,举个例子:

时间 线程 动作
t1 生产者 A CAS 成功,获得 tail=8,对应 slot[0]
t2 消费者 还没消费 slot[0](数据旧)
t3 生产者 B 也 CAS 成功,tail=16,对应 slot[0](再次循环)

        那么,在以上例子,我们是如何知晓slot[0]是已经被消费过了,还是没有被消费,这就需要sequence进行辅助判断。

        所以说sequence就是用来区分这是第几轮循环的slot,当前slot是否可以被再次写入,是否之前存储在这个slot中的数据已经被消费了。

生产者写入逻辑
// Enqueue 往环形缓冲区中放入一个元素
func (r *RingBuffer) Enqueue(v T) bool {
    for {
        // 1️⃣ 读取当前的 tail(写入位置)和 head(读取位置)
        tail := atomic.LoadUint64(&r.tail)
        head := atomic.LoadUint64(&r.head)

        // 2️⃣ 检查是否队列已满
        // 如果 tail - head == capacity,当前tail和head是一直进行增加的,说明环形队列已经占满
        if tail - head >= uint64(len(r.slots)) {
            return false // 满了,不可写
        }

        // 3️⃣ 使用 CAS 尝试“抢占”这个写入序号(类似加锁,但只针对这个数)
        // 多个生产者同时尝试推进 tail,但只有一个会成功。
        if atomic.CompareAndSwapUint64(&r.tail, tail, tail+1) {

            // 4️⃣ 拿到写入位置对应的槽位下标
            idx := tail & r.mask
            slot := &r.slots[idx]

            // 5️⃣ 等待该槽位变为空闲状态
            //   判断条件:slot.sequence == tail,因为tail会一直增加,所以不能将slot.sequence初始化为0
            //   而是应该初始化为从头到当前slot的步数,这样才可以判断是不是走到这步,且处于空闲
            //   如果slot.sequence > tail那就是表明当前
            //   表示该槽目前是空的(旧的数据已被消费完)
            for {
                seq := atomic.LoadUint64(&slot.sequence)
                if seq == tail {
                    break // 说明可写
                }
                runtime.Gosched() // 短暂让出 CPU,避免死自旋
            }

            // 6️⃣ 写入数据(现在这个槽安全地属于当前生产者)
            slot.value = v

            // 7️⃣ 发布数据:让消费者知道这个槽已经可以读了
            //   注意顺序!一定要先写 value,再修改 sequence(发布顺序保证)
            atomic.StoreUint64(&slot.sequence, tail + 1)

            return true
        }

        // 8️⃣ 如果 CAS 失败,说明别的线程抢到了这个 tail
        //    那就回去重试
    }
}

消费者读取逻辑
// Dequeue 从环形缓冲区中取出一个元素
func (r *RingBuffer) Dequeue() (T, bool) {
    for {
        // 1️⃣ 读取 head(当前读取位置)和 tail(生产者最新写入位置)
        head := atomic.LoadUint64(&r.head)
        tail := atomic.LoadUint64(&r.tail)

        // 2️⃣ 检查是否为空:没有可读数据
        if head >= tail {
            return *new(T), false // 空队列
        }

        // 3️⃣ 计算要读取的槽位索引
        idx := head & r.mask
        slot := &r.slots[idx]

        // 4️⃣ 等待该槽位被发布(即生产者已写入)
        //   判断条件:slot.sequence == head + 1
        //   如果还没到,说明生产者还没发布这条数据
        for {
            seq := atomic.LoadUint64(&slot.sequence)
            if seq == head + 1 {
                break // 已发布,可读
            }
            runtime.Gosched()
        }

        // 5️⃣ 读取数据
        v := slot.value

        // 6️⃣ 将该槽位标记为“已消费”,以便下一个生产者可以重用这个位置
        //   这里的 trick:把 sequence 设成 head + capacity
        //   表示“此槽位下一次循环时才会重新可写”
        atomic.StoreUint64(&slot.sequence, head + uint64(len(r.slots)))

        // 7️⃣ 消费者推进 head 指针
        atomic.StoreUint64(&r.head, head + 1)

        return v, true
    }
}

关键点:生产者等待 slot.sequence == tail(或者特定标记),把数据写入并把 slot.sequence 改为 tail+1(已发布),消费者读取时等待 slot.sequence == head+1(已发布),读完把 slot.sequence 改成 head + capacity(腾出槽),表示下一次轮次 / head下一次为多少时候才会再次消费这个数据。

【问题】为什么slot.sequence == tail是判断,那slot.sequence < tail有什么含义,slot.sequence > tail呢?

        因为开始时候进行赋值对应插槽位置是需要附上对应下标值:

slot[0].sequence = 0
slot[1].sequence = 1
slot[2].sequence = 2
...
slot[7].sequence = 7

        其含义是,第0次只能写到sequence为0的插槽位置上第一次只能写到第 1 次写只能写sequence为1的位置上

        那么我们可以看到在消费端进行举例,比如当前消费第8个数据,此时idx = tail & mask = 0tail = 16,消费者在消费后会做slot.sequence = old_tail + capacity = 8 + 8 = 16,所以进行判断:

slot.sequence == 16 == tail

表示当前写入是安全的。所以sequence == tail正好空出来表示等待写入下一轮数据。

        那么slot.sequence < tail有什么含义?这其实是一个非法状态,含义是消费者还没把sequence推入到第一个状态,所以当前此处没有消费写入会覆盖未消费数据,不能写入。

        那slot.sequence > tail有什么含义呢?这种情况其实不应该出现,出现这种情况就是未来轮次到达这个插槽,这通常是bug情况。

【问题】如果出现写入过快,而没有消费者消费,导致当前环形缓冲区全部插槽被占满,无法再写入数据,如何解决?

        这就是典型背压机制,通常会采取的措施有以下四种(流量管控思想):

而最常用的就是阻塞等待,覆盖最久的数据与直接丢弃。

同步与可见性(内存屏障的细节)

首先,理解同步与可见性的问题,就需要理解,为什么CAS成功不等于正确。

【问题】CAS成功为什么不等于正确且成功的操作?

        很多人会认为使用了atomic的CAS操作就会安全,说明线程是绝对安全的,不会有并发问题出现。但是这个是错误的。

        因为CAS仅仅只能保证某个变量的修改是原子性,而不保障修改的顺序和可见性。采用CAS确实可以保障获取写入权,但是不能保障“其他CPU在什么时候可以看到你写入的内容”。尤其在无锁的结构中,顺序与可见性比原子性更为关键。

【注意】典型的发布-订阅顺序错误问题:

        在环形缓冲区+CAS的结构中,sequence就会存在对应问题,假设生产者的操作被CPU重新排序,则会出现:

atomic.StoreUint64(&slot.sequence, tail+1) // 发布先执行
slot.value = v                              // 数据后写

那就会导致消费者看到sequence已经被迭代了,代表当前插槽内存储的是新数据,但是其实目前插槽内新数据没有被写入,如果读取则读到的是脏数据(未写完的旧数据)。

        正确的发布-订阅模型的执行顺序是:

操作主体 步骤 操作 说明
生产者 写入数据 → slot.value = v 写业务数据
生产者 原子发布 → atomic.StoreUint64(&slot.sequence, tail+1) 发布标记(Release)
消费者 检查发布标记 → atomic.LoadUint64(&slot.sequence) 判断是否可读(Acquire)
消费者 读取数据 → slot.value 获取已发布数据

所以,为了保障发布-订阅执行顺序,CPU采用了内存屏障。

        在多核CPU的情况下,每个核都有自己的缓存层(通常有两层,L1/L2),所以即使一个线程修改了内存中的某个值,对应CPU的缓存以及共享主内存也没有及时更新,而其他CPU读取当前被修改后的数据,是先去当前共享主内存中拿取,这就会导致拿到的缓存内容其实跟当前CPU内部存储的数据是不一样的。(总结:强制所有读取内存中数据的操作,都从CPU中读取,不从缓存中,且会将CPU中数据同步更新到主内存与缓存

        内存屏障模型相关指令与含义:

而Go中具体的方法与CPU指令的实现:

ABA问题和防护

        ABA问题:就是线程读取内存数据是A,然后进行CAS操作,在这个过程中有其他线程将这个线程数据改为B,但是又来了一个线程将数据B改为A,这样在当前线程计算完毕最终数据后进行oldValue和当前位置数据时候,是一致的,但是其实并不是之前的数据,这就是ABA问题。CAS只能识别数据是否相等,无法识别“中间改变过与否”。

        而环形缓冲区中就使用sequence即插槽版本号进行解决这个问题。

环形缓冲区+CAS的优化

为了更高吞吐和减少 contention,实施中常做:

  • Batch produce / batch consume:生产者一次写多个事件并一起 publish(减少 tail 更新次数)。消费者也一次拉多条处理。

  • 背压(backpressure):当 buffer 快满时,使用返回 false / block /缓冲排队 等策略通知上游降速。

  • 选择合理 capacity:容量太小容易溢出/高竞争;太大内存浪费且 cache 局部性下降。通常根据吞吐量和延迟实验确定。

常见陷阱与具体实现或者情况展示/解决思路
  • 忘记发布序列前写数据 → 消费者读取到空/半写数据:总是先写数据,再做 atomic publish(release)。

  • head/tail 未做 padding → false sharing:CPU 缓存行频繁失效,吞吐暴跌。

  • 未处理 wrap-around/overflow:使用 64-bit 序号并设计检测。

  • 使用短整数做 seq(如 32-bit)在高吞吐下容易回绕

  • 频繁 GC 的 value 导致性能不稳:预分配或使用对象池。

  • 过度自旋造成 CPU 饱和:实现退避并在长期延迟时切换到阻塞。

  • 多生产者直接对 slot.value CAS(错误):不要在 slot.value 上 CAS,应该先 claim seq 再写。

Logo

汇聚全球AI编程工具,助力开发者即刻编程。

更多推荐