Skip to content

环形缓冲区(Ring Buffer) ​

#数据结构 · #环形缓冲区 · #队列 · #无锁 · #Kafka · #Go Channel

环形缓冲区(Ring Buffer / Circular Buffer)是一种固定大小的 FIFO 队列,用两个指针在数组上循环移动。它在生产者-消费者场景中性能极高,广泛应用于消息队列、网络 I/O、音视频流等领域。


基本原理 ​

结构 ​

数组大小 = 8
  write                              read
    ↓                                  ↓
  ┌───┬───┬───┬───┬───┬───┬───┬───┐
  │ 5 │ 6 │ 7 │   │   │ 1 │ 2 │ 3 │
  └───┴───┴───┴───┴───┴───┴───┴───┘
    ↑                              ↑
  末尾                            起始

可写空间 = (read - write - 1 + N) % N
可读元素数 = (write - read + N) % N

核心规则:

  • write 指向下一个写入位置
  • read 指向下一个读取位置
  • 始终保持至少一个空位,以区分"满"和"空"

基本操作 ​

go
type RingBuffer struct {
    buf   []interface{}
    size  int
    read  int
    write int
}

func NewRingBuffer(capacity int) *RingBuffer {
    return &RingBuffer{
        buf:  make([]interface{}, capacity+1), // +1 用于区分满/空
        size: capacity + 1,
    }
}

func (r *RingBuffer) IsEmpty() bool {
    return r.read == r.write
}

func (r *RingBuffer) IsFull() bool {
    return (r.write+1)%r.size == r.read
}

func (r *RingBuffer) Write(val interface{}) bool {
    if r.IsFull() {
        return false
    }
    r.buf[r.write] = val
    r.write = (r.write + 1) % r.size
    return true
}

func (r *RingBuffer) Read() (interface{}, bool) {
    if r.IsEmpty() {
        return nil, false
    }
    val := r.buf[r.read]
    r.buf[r.read] = nil  // 帮助 GC
    r.read = (r.read + 1) % r.size
    return val, true
}

func (r *RingBuffer) Len() int {
    return (r.write - r.read + r.size) % r.size
}

func (r *RingBuffer) Cap() int {
    return r.size - 1
}

并发安全:无锁环形缓冲区 ​

单生产者 / 单消费者(SPSC) ​

go
// 利用内存屏障保证可见性(简化示意,实际需 atomic)
type LockFreeRingBuffer struct {
    buf        []interface{}
    mask       int    // size 必须是 2 的幂,mask = size - 1
    readIdx    int64  // atomic
    writeIdx   int64  // atomic
    readCache  int64  // 消费者缓存的 writeIdx
    writeCache int64  // 生产者缓存的 readIdx
}

func (b *LockFreeRingBuffer) Write(val interface{}) bool {
    write := atomic.LoadInt64(&b.writeIdx)
    if write-b.writeCache >= int64(b.mask+1) {
        b.writeCache = atomic.LoadInt64(&b.readIdx)
        if write-b.writeCache >= int64(b.mask+1) {
            return false  // 满了
        }
    }
    b.buf[write&int64(b.mask)] = val
    atomic.StoreInt64(&b.writeIdx, write+1)
    return true
}

func (b *LockFreeRingBuffer) Read() (interface{}, bool) {
    read := atomic.LoadInt64(&b.readIdx)
    if read >= b.readCache {
        b.readCache = atomic.LoadInt64(&b.writeIdx)
        if read >= b.readCache {
            return nil, false  // 空了
        }
    }
    val := b.buf[read&int64(b.mask)]
    atomic.StoreInt64(&b.readIdx, read+1)
    return val, true
}

多生产者 / 多消费者(MPMC) ​

需要 CAS 同时竞争 read/write 槽位。Disruptor(LMAX 开源的超高性能队列)采用了更彻底的方案:

mermaid
flowchart TB
    subgraph Disruptor["Disruptor (LMAX)"]
        direction TB
        RB["RingBuffer<br/>预分配固定大小数组<br/>所有数据在一开始就分配好"]

        subgraph Producers["生产者"]
            P1["Producer 1<br/>CAS 申请下一个槽位<br/>写入数据<br/>更新 Sequence"]
            P2["Producer 2"]
        end

        subgraph Consumers["消费者"]
            C1["Consumer 1<br/>通过 Sequence Barrier<br/>等待 Sequence 到达<br/>批量读取处理"]
            C2["Consumer 2"]
        end

        Producers -->|"CAS 竞争 cursor"| RB
        RB -->|"Sequence Barrier"| Consumers
    end

Disruptor 的核心创新:

  1. 预分配所有数据:RingBuffer 在创建时就分配好所有槽位对象,运行时零 GC
  2. 缓存行填充:Sequence 值前后各填充 7 个 long(64B 缓存行),消除伪共享
  3. 批量消费:Consumer 一次处理一批 Sequence,减少 CAS 操作
  4. 无锁全程:仅用 CAS 申请槽位,用内存屏障(volatile/Ordered)保证可见性
go
// Go 版 Disruptor 核心思想(简化)
type Sequence struct {
    value   int64
    padding [56]byte  // 缓存行填充:防止伪共享
}

// 伪共享 (False Sharing):
// CPU 缓存行 = 64B
// 两个 Core 分别更新相邻 8B 的数据 →
//   虽然逻辑无关,但缓存一致性协议(MESI)会不断使对方缓存失效
//   实际性能损失可达 10-100x
//
// Disruptor 通过对 Sequence 做 64B 对齐来避免这个问题

Linux kfifo — 内核中的 Ring Buffer ​

kfifo 是 Linux 内核中最简洁高效的环形缓冲区实现,核心设计就是两个 unsigned int 的自然溢出:

c
// kfifo 核心操作(无需锁,因为 in/out 分别由生产者和消费者独占)
struct __kfifo {
    unsigned int in;    // 生产者维护,自然溢出回绕
    unsigned int out;   // 消费者维护
    unsigned int mask;  // size - 1(2 的幂)
    char data[];        // 柔性数组
};

// in 和 out 用 unsigned int,溢出后自动回绕到 0
// 假设 in = 0xFFFFFFFF, in++ → 0x00000000(这正是我们需要的!)
特性kfifoDisruptorGo Channel
并发模型无锁 SPSCCAS MPMCMutex(有缓冲时环形)
阻塞/非阻塞非阻塞非阻塞阻塞
批量操作✅ 一次 memcpy 可写满✅ BatchConsumer❌ 逐元素
缓存行优化❌(内核场景不重要)✅ 核心优化❌
使用场景内核驱动、硬件间数据传递金融交易(单机百万 TPS)goroutine 间通信

各系统中的实现 ​

1. Go Channel 底层 ​

go
// hchan 结构中的环形缓冲区(简化)
type hchan struct {
    buf      unsafe.Pointer  // 环形缓冲区
    elemsize uint16
    qcount   uint            // 当前元素数
    dataqsiz uint            // 缓冲区大小
    sendx    uint            // 发送索引(write)
    recvx    uint            // 接收索引(read)
}
mermaid
flowchart LR
    subgraph Channel["Go Channel (buf=3)"]
        direction TB
        BUF["[A][B][ ]"]
        SENDX["sendx=2"]
        RECVX["recvx=0"]
    end

    G1["goroutine-1<br/>ch <- 'C'"] -->|写入 buf[2]| BUF
    BUF -->|sendx 变为 0<br/>(2+1)%3=0| SENDX

    G2["goroutine-2<br/>v := <-ch"] -->|读取 buf[0]='A'| BUF
    BUF -->|recvx 变为 1| RECVX

2. Linux kfifo ​

c
// linux/kfifo.h — 内核中的环形缓冲区
struct __kfifo {
    unsigned int in;    // write 偏移
    unsigned int out;   // read 偏移
    unsigned int mask;  // size - 1(必须 2 的幂)
    // data 跟在结构体后面
};

// 写入(简化)
unsigned int kfifo_in(struct __kfifo *fifo, const void *buf, unsigned int len) {
    unsigned int l;
    len = min(len, fifo->mask + 1 - fifo->in + fifo->out);
    // 第一部分:从 in 到末尾
    l = min(len, fifo->mask + 1 - (fifo->in & fifo->mask));
    memcpy(fifo->data + (fifo->in & fifo->mask), buf, l);
    // 第二部分(可能回绕):从开头
    memcpy(fifo->data, buf + l, len - l);
    fifo->in += len;
    return len;
}

关键优化:in 和 out 用 unsigned int 自然溢出,除以 2 的幂用 & mask 代替 %。

3. Kafka 日志段 ​

Kafka Partition 的日志存储:
┌─────────────────────────────────────────────────────┐
│ Partition                                            │
│  offset →  0  1  2  3  4  5  6  7  8  9  10  ...   │
│  ─────────────────────────────────────────────────── │
│  Segment 1 (0-4)     │ Segment 2 (5-9)  │ ...       │
│  [0][1][2][3][4]     │ [5][6][7][8][9]  │           │
└─────────────────────────────────────────────────────┘

虽然不是严格环形,但思想类似:
- 写入顺序追加(append-only)
- 旧数据按时间或大小淘汰

4. 音视频流处理 ​

PCM 音频采集 → Ring Buffer → 编码器 → Ring Buffer → 网络发送
              ↑ (读写同时进行,零拷贝)

环形缓冲区 vs 普通队列 ​

维度环形缓冲区链表队列动态数组队列
内存分配一次分配每次入队扩容时分配
缓存友好✅ 连续内存❌ 指针跳转✅ 连续内存
内存占用固定,始终占用按需增长可缩可扩
GC 压力低(无对象分配)高(每节点一个对象)中
无锁实现难度易(SPSC)难中
适用场景容量已知、高频读写容量不可预测容量基本可预测

工程实践:背压、局部性与热点 ​

1. Ring Buffer 满了以后怎么办? ​

这是工程上最重要的问题之一。算法课里通常只说“满了返回失败”,但真实系统必须明确背压策略。

策略做法优点风险适用场景
阻塞等待满了就等消费者腾空间不丢数据上游线程可能被拖死任务必须保序且不能丢
直接丢弃返回失败,调用方重试/丢弃简单、保护系统数据可能丢失日志、指标、采样
覆盖旧数据新数据覆盖最旧槽位始终有最新数据历史数据丢失监控、实时看板
扩容迁移分配更大数组再搬迁容量弹性失去 ring buffer 固定内存优势普通队列更常见
反压上游降速、限流、拒绝全链路稳定实现复杂高吞吐服务、消息系统

很多系统最后不是败在 ring buffer 本身,而是没定义“满了以后谁承担代价”。

2. 为什么它对 CPU 友好 ​

Ring Buffer 在工程上非常受欢迎,核心不是时间复杂度,而是局部性(locality)。

text
数组连续内存
  → 预取友好
  → cache line 命中率高
  → TLB miss 少
  → 比链表队列更适合高频读写

这也是为什么很多高性能系统宁愿用固定长度数组 + 序号,而不是链表:

  • 链表理论上插入删除也很快
  • 但节点分散在堆上,cache miss 很重
  • 真正拖慢系统的往往不是 O(1) / O(log n),而是访存和缓存一致性开销

3. SPSC / MPSC / MPMC 的难点差别 ​

模型难点常见实现
SPSC最简单,读写各自独占索引原子读写即可
MPSC多个生产者竞争写指针CAS 竞争 cursor
SPMC多个消费者竞争读指针CAS 竞争 read index
MPMC读写双方都竞争每个槽位 sequence / ticket 最常见

复杂度真正上升的原因不是“多几个线程”,而是:

  • CAS 失败重试
  • 内存屏障增多
  • 伪共享导致 cache line ping-pong
  • 空转自旋浪费 CPU

4. CAS 热点与伪共享 ​

在高并发场景里,Ring Buffer 的瓶颈经常集中在少数几个共享变量上:

  • 全局 writeIdx
  • 全局 readIdx
  • cursor / sequence

如果多个核心频繁更新同一 cache line,会出现:

text
Core 1 写 cursor
  → Core 2 的缓存行失效
Core 2 再写
  → Core 1 的缓存行又失效

于是你看到的现象可能是:

  • 算法上仍然是 O(1)
  • 但吞吐上不去
  • CPU 很高,实际有效工作不多

这也是 Disruptor 要做缓存行填充的根本原因。


实战案例 ​

1. 为什么日志队列容易丢数据 ​

很多异步日志库底层就是 ring buffer。高峰期磁盘或网络写出速度不够时,如果采用“队列满就丢”,表现就是:

  • 业务线程几乎不受影响
  • 但日志出现缺口
  • 故障时偏偏缺关键日志

所以日志系统必须明确:

  • 是允许丢日志,还是宁可拖慢业务也要保留
  • 是按级别丢弃,还是按采样率丢弃
  • 是否要暴露 dropped count 指标

2. 为什么交易/撮合系统偏爱 Ring Buffer ​

在金融撮合、游戏事件循环、音视频流水线里,数据路径通常是固定大小、低延迟、强局部性,这正好适合 ring buffer:

  • 内存预分配,避免运行时分配
  • 批量处理,减少 CAS 和系统调用
  • 访问模式稳定,缓存命中高

3. 什么时候不该用 Ring Buffer ​

以下场景反而不适合:

  • 数据量高度不可预测
  • 不能接受丢弃,但也不能阻塞上游
  • 需要复杂优先级调度
  • 消费速度和生产速度长期不匹配

这时更适合:

  • 带持久化的消息队列
  • 支持动态扩容的普通队列
  • 带优先级的堆或多队列模型

面试要点 ​

经典问题 ​

  1. 区分满和空:留一个空位 OR 加 count 字段

    空: read == write
    满: (write + 1) % size == read (留一个空位)
    或: count == size (额外字段)
  2. 为什么大小用 2 的幂:% size 换成 & (size-1),快一个数量级

  3. 如何保证并发安全:

    • SPSC:无锁,依赖原子操作 + 缓存行填充
    • MPMC:CAS 竞争槽位,Disruptor
  4. 与阻塞队列的区别:

    • 环形缓冲区满了返回 false,由调用方处理
    • 阻塞队列满了会 block,依赖条件变量唤醒

工业实现对比:kfifo vs Go channel vs Kafka ​

1. 三种"环形思想"的异同 ​

维度Linux kfifoGo channel (buf>0)Kafka log segment
数据结构环形缓冲区(数组)环形缓冲区(数组)顺序追加文件(非环形)
内存管理固定大小,预分配固定 buffer 大小文件系统管理,可无限扩展
并发模型无锁 SPSC/MPMC有锁(mutex)单线程 append + 分区隔离
满时行为覆盖或返回错误阻塞(send)写下一个 segment
消费方式批量读取单元素 pop按 offset 随机读
适用场景内核数据路径goroutine 间通信消息持久化 + 多 consumer
Cache 友好性✅ 连续数组✅ 连续数组🟡 顺序文件 IO

2. 为什么 kfifo 用无锁,Go channel 用有锁? ​

text
kfifo (Linux kernel):
  - C 语言,直接操作内存
  - 编译时知道是 SPSC 还是 MPMC
  - 无锁靠: 原子变量 + 内存屏障 (smp_mb / smp_wmb)
  - 原因: 内核不能容忍锁带来的延迟不确定性

Go channel:
  - Go 语言的 channel 是通用抽象
  - 编译时不知道是单生产者还是多生产者
  - 有锁靠: runtime.mutex (futex based)
  - 原因: Go 的 M:N 调度模型中,mutex 阻塞时会让出 P
          → 不浪费 CPU → 锁的开销可接受

3. SPSC vs MPMC 的内存屏障差异 ​

c
// SPSC (单生产者单消费者) — 最少屏障
// 生产者只需写屏障,消费者只需读屏障

// 生产者: store(write_idx) 前确保数据可见
data[write_idx] = val;
smp_wmb();  // 写屏障: 保证 data 写完再更新 write_idx
WRITE_ONCE(ring->write_idx, next_write);

// 消费者: load(write_idx) 后确保数据可见
smp_rmb();  // 读屏障: 保证读到 write_idx 后再读 data
val = READ_ONCE(ring->data[read_idx]);

// MPMC (多生产者多消费者) — 需要 CAS + 全屏障
// 生产者需要 CAS 竞争 write_idx
do {
    cur = ring->write_idx;
    next = cur + 1;
} while (!CAS(&ring->write_idx, cur, next));
// 全屏障: 需要 smp_mb() 保证 CAS 可见

4. Disruptor 的关键设计:缓存行填充避免伪共享 ​

Disruptor 的性能核心之一是缓存行填充(Cache Line Padding):

java
// Disruptor Sequence 的伪共享防护
public class Sequence {
    // 前 7 个 long (56B) + value (8B) = 64B = 一个 cache line
    protected long p1, p2, p3, p4, p5, p6, p7;
    private volatile long value;
    protected long p8, p9, p10, p11, p12, p13, p14;
}
text
没有 padding:
  Producer 的 cursor 和 Consumer 的 cursor 在同一 cache line
  → Producer 写 cursor → Consumer 的 cache line 失效
  → Consumer 读 cursor → Cache miss → 跨核通信
  → 惨案: 本来 1ns 的 L1 操作变成 ~100ns 的跨核通信

有 padding:
  两个 cursor 各自独占一个 cache line
  → 互不干扰 → 无伪共享 → 真正的无锁性能

5. Disruptor 为什么快? ​

优化原理收益
Ring Buffer 预分配无运行时分配消除 GC 压力
缓存行填充避免 false sharing消除跨核 cache 失效
无锁 CAS避免锁竞争高并发下线性扩展
批量提交生产者一次写多条,一次 commit减少 CAS 次数
Event Processor 依赖图多个 consumer 可以并行处理不同阶段流水线化

6. 选型指南 ​

场景推荐原因
goroutine 间通信Go channel语言原生,有 select/mutex 配合
内核高速数据路径Linux kfifo无锁、零拷贝、确定性延迟
微秒级交易撮合Disruptor极致低延迟 + 批量处理
消息持久化 + 多 consumerKafka磁盘持久化 + 消费者组隔离
日志缓冲io.Writer + ring buffer批量 flush + 背压控制
音频/视频流水线环形缓冲区固定帧率 + DMA 友好

参考 ​

批注模式

💬 文章评论

暂无评论,来说点什么吧 👇

编程学习笔记