# pipeline_ringbuffer **Repository Path**: deaglebear/pipeline_ringbuffer ## Basic Information - **Project Name**: pipeline_ringbuffer - **Description**: pipeline_ringbuffer - **Primary Language**: Unknown - **License**: Apache-2.0 - **Default Branch**: master - **Homepage**: None - **GVP Project**: No ## Statistics - **Stars**: 0 - **Forks**: 0 - **Created**: 2025-11-09 - **Last Updated**: 2026-07-30 ## Categories & Tags **Categories**: Uncategorized **Tags**: None ## README # Pipeline RingBuffer [![License: MIT](https://img.shields.io/badge/License-MIT-yellow.svg)](LICENSE) [![C++ Standard](https://img.shields.io/badge/C%2B%2B-20-blue.svg)](https://en.cppreference.com/w/cpp/20) 一个高性能、基于原子操作的多阶段流水线环形缓冲区,使用 C++20 实现。适用于数据需要流经多个处理阶段的场景,提供精细的并发控制和伪共享防护。 ## 目录 - [概述](#概述) - [适用场景](#适用场景) - [核心概念](#核心概念) - [编译与安装](#编译与安装) - [快速入门](#快速入门) - [Step 类型详解](#step-类型详解) - [ExclusiveStep(独占步骤)](#exclusivestep独占步骤) - [ShareStep(共享步骤)](#sharestep共享步骤) - [Consumer API](#consumer-api) - [进阶用法](#进阶用法) - [伪共享防护](#伪共享防护) - [线程安全](#线程安全) - [API 参考](#api-参考) - [测试](#测试) --- ## 概述 `PipelineRingBuffer` 实现了一个**跨多个处理阶段共享的固定大小环形缓冲区**。缓冲区中的每个元素携带用户数据(`T`)以及一个内部状态机,用于控制哪个阶段可以获取该元素。 核心特性: - **多阶段流水线**:数据自动按照 阶段 0 → 阶段 1 → … → 阶段 N → 回到阶段 0 的顺序流转 - **两种步骤类型**:Exclusive(每个元素由一个消费者独占,基于 CAS 抢占)和 Share(多消费者,基于提交计数协同) - **伪共享防护**:每个缓冲区槽位都有填充,避免缓存行竞争 - **纯头文件**:`#include` 即可使用,无需编译链接 - **无互斥锁热路径**:领取和提交只使用原子操作,但不承诺严格的 lock-free 进度保证 ## 适用场景 本库适用于以下场景: - 需要**按阶段顺序处理数据流**(如:解析 → 转换 → 校验 → 输出) - 每个阶段运行在独立的线程中 - 需要**限制内存使用**(固定缓冲区大小) - 需要**低延迟、高吞吐**的阶段间数据传递 - 多个消费者需要**并行协作处理**同一阶段的数据(ShareStep) 如果只需要简单的单生产者-单消费者队列,使用 `std::queue` + 互斥锁可能更简单。本库的优势在于**流水线拓扑结构**——多个阶段串联处理。 ## 核心概念 ``` 阶段 0 阶段 1 阶段 2 阶段 3 (Exclusive) (Exclusive) (Share) (Exclusive) │ │ │ │ ┌───▼───┐ ┌───▼───┐ ┌───▼───┐ ┌───▼───┐ │Consumer│ │Consumer│ │Consumer│ │Consumer│ │ (输入) │ ──▶ │ (解析) │ ──▶ │ (处理) │ ──▶ │ (输出) │ └───────┘ └───────┘ │Consumer│ └───────┘ │Consumer│ └───────┘ ▲ 多个消费者并行处理同一阶段 ``` **RingBuffer(环形缓冲区)** — 共享容器。它持有一个 `Item` 数组和一个 `Step` 对象列表。 **Step(步骤)** — 表示流水线中的一个处理阶段。每个步骤有类型(Exclusive 或 Share)和一组消费者。 **Consumer(消费者)** — 工作线程用于从所在步骤获取元素、处理元素、并将元素提交到下一步骤的句柄。 **Item(元素)** — 缓冲区中的一个槽位。包含用户数据(`T`)和内部状态字段(`status`、`commit_count`、`index`)。 数据在流水线中的流转过程: ``` [阶段 0 可用] → Consumer 0 获取 → 处理中 → Consumer 0 提交 ↓ [阶段 1 可用] → Consumer 1 获取 → 处理中 → Consumer 1 提交 ↓ … 更多阶段 … ↓ [阶段 0 可用] ← ← ← ← ← ← ← ← ← ← ← ← ← ← ← ← ┘ ``` ## 编译与安装 ### 环境要求 - **C++20** 编译器(GCC 11+、Clang 16+) - **CMake** ≥ 3.15 ### 编译项目 ```bash git clone <仓库地址> cd pipeline_ringbuffer cmake -S . -B build cmake --build build ``` ### 在项目中使用(CMake) ```cmake # 方式一:add_subdirectory add_subdirectory(path/to/pipeline_ringbuffer) target_link_libraries(your_target PRIVATE pipeline_ringbuffer) # 方式二:install + find_package # (先执行 cmake --install build --prefix /your/prefix) find_package(pipeline_ringbuffer REQUIRED) target_link_libraries(your_target PRIVATE pipeline_ringbuffer) ``` 本库是**纯头文件库**。`pipeline_ringbuffer` 这个 CMake target 会提供正确的头文件路径和 C++20 编译特性。 ### 依赖项 核心库**无外部依赖**——仅使用 C++ 标准库。测试代码使用了 [rapidjson](https://github.com/Tencent/rapidjson)、[spdlog](https://github.com/gabime/spdlog) 和 [magic_enum](https://github.com/Neargye/magic_enum) 用于基准测试和日志(均为可选依赖)。 ## 快速入门 最简单的流水线:一个输入阶段 + 一个输出阶段。 ```cpp #include "pipeline_ringbuffer/pipeline_ringbuffer.h" #include #include using namespace PipelineRingBuffer; struct MyData { int value; }; int main() { // 创建环形缓冲区:2 个阶段,类型为 Exclusive → Exclusive // buffer_size 必须是 2 的幂 RingBuffer ring_buffer(8, {EStepType::Exclusive, EStepType::Exclusive}); // 为每个阶段创建消费者 auto producer = ring_buffer.create_consumer(0); // 阶段 0 auto consumer = ring_buffer.create_consumer(1); // 阶段 1 // 生产者线程:写入数据 std::thread producer_thread([&]() { for (int i = 0; i < 100; i++) { auto item = producer->claim_guard(); // 阻塞直到有可用元素 item->data.value = i; } // claim_guard() 返回的 shared_ptr 析构时自动提交 }); // 消费者线程:读取数据 std::thread consumer_thread([&]() { for (int i = 0; i < 100; i++) { auto item = consumer->claim_guard(); std::cout << "收到: " << item->data.value << std::endl; } }); producer_thread.join(); consumer_thread.join(); return 0; } ``` ## Step 类型详解 ### ExclusiveStep(独占步骤) **独占**步骤保证每个元素**只有一个消费者**能够获取。多个消费者可以绑定到同一个独占步骤,但它们通过 **CAS 竞争**——每个元素只有一个消费者能抢到。 领取过程: 1. 读取单调递增的 64 位 `m_sequence`,通过 `sequence & cap_mask` 计算槽位索引。 2. 使用 CAS 将槽位 `status` 从当前阶段的 available 原子地转换为 unavailable,获得该槽位的唯一所有权。 3. 使用 CAS 推进 `m_sequence`。如果游标已被其他线程推进,则恢复槽位状态并重试。 4. `commit()` 使用 release store 将 `status` 发布为下一阶段的 available。 `m_sequence` 同时承担游标和轮次编号的职责。完整 sequence 参与 CAS,可以阻止旧游标值在绕环后错误匹配;槽位索引则由低位推导,不需要额外维护 `{value, term}` 组成的 16 字节原子对象。在常见的 64 位平台上,`std::atomic` 通常可以由原生指令实现,具体平台可通过 `std::atomic::is_always_lock_free` 检查。 该实现保证不会有两个 ExclusiveConsumer 同时获得同一槽位。不过,如果线程在成功占用 `status` 后、推进 `m_sequence` 前暂停,当前阶段会暂时停止推进,因此这里的“无互斥锁”不等同于严格的 lock-free。 适用场景: - 每个阶段只有一个工作线程 - 少量工作线程竞争处理任务 ```cpp // 三个工作线程都从同一个 ExclusiveStep 消费 auto w1 = ring_buffer.create_consumer(1); auto w2 = ring_buffer.create_consumer(1); auto w3 = ring_buffer.create_consumer(1); // 它们通过 CAS 竞争,每个元素只会被一个工作线程获取 ``` ### ShareStep(共享步骤) **共享**步骤允许多个消费者**同时处理同一个元素**。每个消费者有自己独立的游标(`m_tail`),因此它们独立遍历缓冲区。元素只有在**所有消费者都提交后**才会推进到下一阶段。 适用场景: - 需要对同一个元素进行并行处理(如:多个线程并发解析同一个 JSON 的不同字段) - 需要**栅栏语义**——等待所有工作线程完成后再进入下一阶段 ```cpp // 三个工作线程都处理相同的元素 auto w1 = ring_buffer.create_consumer(2); // ShareStep auto w2 = ring_buffer.create_consumer(2); auto w3 = ring_buffer.create_consumer(2); // 三个工作线程都能看到每个元素;只有全部提交后元素才会进入下一阶段 ``` **注意事项**:ShareStep 的消费者各自独立跟踪位置。它们必须以大致相同的速度消费元素——**最慢的消费者决定了整体吞吐量**。 ## Consumer API 每个消费者提供以下几种获取元素的策略: | 方法 | 行为 | 返回值 | |------|------|--------| | `try_claim()` | 非阻塞:无可用元素时立即返回 | `Item*` 或 `nullptr` | | `claim()` | 阻塞:自旋直到有可用元素 | `Item*` | | `claim(timeout)` | 超时阻塞:超时后返回 `nullptr` | `Item*` 或 `nullptr` | | `try_claim_guard()` | 非阻塞 RAII:析构时自动提交 | `shared_ptr>` 或 `nullptr` | | `claim_guard()` | 阻塞 RAII:自旋直到可用,自动提交 | `shared_ptr>` | | `claim_guard(timeout)` | 超时 RAII:超时返回 `nullptr` | `shared_ptr>` 或 `nullptr` | | `commit(item)` | 显式标记元素处理完成(仅 `try_claim` / `claim` 需要手动调用) | void | **建议**:尽量使用 `_guard` 版本——它们会在 `shared_ptr` 离开作用域时自动调用 `commit()`,避免忘记提交导致流水线阻塞。 ```cpp // RAII 风格 —— 推荐 { auto item = consumer->claim_guard(); // 阻塞直到可用 item->data.do_work(); } // 离开作用域,自动提交 // 手动风格 —— 需要精细控制时使用 Item* item = consumer->try_claim(); if (item) { item->data.do_work(); consumer->commit(item); // 千万别忘了这一步! } ``` ## 进阶用法 ### 混合步骤类型的多阶段流水线 以下示例展示一个 4 阶段流水线:输入 → 字段切分(并行)→ 数据解析(并行)→ 输出。 ```cpp #include "pipeline_ringbuffer/pipeline_ringbuffer.h" using namespace PipelineRingBuffer; struct QuoteItem { std::string raw_json; double price; double volume; // ... 其他字段 }; // 流水线布局:Exclusive → Exclusive → Share → Exclusive RingBuffer ring_buffer(16, { EStepType::Exclusive, // 阶段 0: 输入(1 个 worker) EStepType::Exclusive, // 阶段 1: 预处理(3 个 worker 竞争) EStepType::Share, // 阶段 2: 并行解析(5 个 worker,每个都处理同一元素) EStepType::Exclusive // 阶段 3: 输出(1 个 worker) }); // 阶段 0: 单个生产者 auto input = ring_buffer.create_consumer(0); // 阶段 1: 多个 worker 竞争获取元素 auto preproc_1 = ring_buffer.create_consumer(1); auto preproc_2 = ring_buffer.create_consumer(1); auto preproc_3 = ring_buffer.create_consumer(1); // 阶段 2: 并行处理 —— 所有 worker 都看到每个元素 auto parser_1 = ring_buffer.create_consumer(2); auto parser_2 = ring_buffer.create_consumer(2); auto parser_3 = ring_buffer.create_consumer(2); // 阶段 3: 单个消费者 auto output = ring_buffer.create_consumer(3); ``` ### 配合停止标志使用 try_claim 当需要优雅地停止工作线程时: ```cpp std::atomic running{false}; void worker_thread(Consumer* consumer) { while (!running) { /* 自旋等待启动信号 */ } while (running) { auto item = consumer->try_claim(); if (!item) continue; // 当前无可用元素,重试 // 处理元素 ... consumer->commit(item); } } // 在主线程中: running = true; std::this_thread::sleep_for(std::chrono::seconds(5)); running = false; // 所有 worker 退出循环 ``` ## 伪共享防护 当多个 CPU 核心访问恰好共享同一缓存行的相邻内存位置时,性能会因**伪共享**(false sharing)而下降。本库通过两种方式防止伪共享: ### 1. 元素级填充 每个 `Item` 开头有 **64 字节的填充**(`char[64]`): ``` ┌───────────── 元素[0] ─────────────┬───────────── 元素[1] ─────────────┐ │ 64B 填充 │ index │ status │ … │ 64B 填充 │ index │ status │ … │ └───────────────────────────────────┴───────────────────────────────────┘ 缓存行边界 缓存行边界 ``` 元素[N+1] 开头的填充充当了元素[N] 的数据部分与元素[N+1] 的原子状态字段之间的缓冲区。这样,当线程 A 写入元素[0].data 而线程 B 同时读取元素[1].status 时,它们不会争用同一条缓存行。 ### 2. 步骤级对齐 `ExclusiveStep::m_sequence` 和 `ShareStep::m_header` 原子游标都按 `hardware_destructive_interference_size`(x86-64 上为 64 字节)对齐,防止步骤内部状态与相邻数据发生伪共享。 填充大小在编译期确定: - 支持 `__cpp_lib_hardware_interference_size` 的编译器:使用平台推荐值 - 降级方案:**64 字节**(x86-64 上正确;ARM/Apple Silicon 的 128B 缓存行场景下偏保守) ## 线程安全 - **先创建拓扑再启动线程**:应在工作线程启动前创建全部 Consumer;处理期间不要并发调用 `create_consumer()` - **ExclusiveStep**:多个 ExclusiveConsumer 可以并发竞争同一个步骤;`status` CAS 保证每个槽位只有一个 owner - **ShareStep**:每个工作线程应使用独立的 ShareConsumer;不要在多个线程间共享同一个 ShareConsumer,因为它包含线程私有的非原子游标 - **元素数据访问不做通用同步**:ShareStep 的多个消费者同时修改 `item->data` 时,需要自行使用互斥锁、原子变量或互不重叠的字段 - **生命周期**:RingBuffer 必须晚于其 Consumer 和尚未析构的 claim guard 销毁 ## API 参考 ### `RingBuffer` ```cpp RingBuffer(uint64_t buffer_size, const std::vector& step_type_list) ``` 构造环形缓冲区。`buffer_size` 必须是 2 的正整数次幂。`step_type_list` 定义流水线中各阶段的类型。 ```cpp Consumer* create_consumer(int step_no) ``` 创建绑定到指定阶段索引的消费者。消费者由 RingBuffer 持有,**不要手动 delete**。 ```cpp const Item* peer(uint64_t index) const ``` 按索引查看任意缓冲区槽位(调试用)。 ### `Consumer`(基类) | 方法 | 说明 | |------|------| | `Item* try_claim()` | 非阻塞获取 | | `Item* claim()` | 阻塞获取(自旋) | | `Item* claim(milliseconds)` | 带超时的获取,超时返回 `nullptr` | | `shared_ptr> try_claim_guard()` | 非阻塞 RAII 获取 | | `shared_ptr> claim_guard()` | 阻塞 RAII 获取 | | `shared_ptr> claim_guard(milliseconds)` | 带超时的 RAII 获取 | | `void commit(Item*)` | 将元素释放到下一阶段 | | `Step* step()` | 获取此消费者所属的步骤 | | `IRingBuffer* ringbuffer()` | 获取父级环形缓冲区 | ### `Item` | 成员 | 类型 | 说明 | |------|------|------| | `data` | `T` | 用户自定义数据 | | `index` | `uint64_t` | 槽位在缓冲区中的索引(只读) | ### `EStepType` | 值 | 含义 | |----|------| | `Exclusive` | 每个元素只有一个消费者获取(基于 CAS) | | `Share` | 多个消费者共享;所有消费者都能看到每个元素 | ### `Step`(基类) | 方法 | 说明 | |------|------| | `step_type()` | 返回 `EStepType::Exclusive` 或 `EStepType::Share` | | `step_no()` | 此步骤在流水线中的索引 | | `ringbuffer()` | 父级环形缓冲区 | | `consumer_list()` | 绑定到此步骤的所有消费者 | ### 填充工具 ```cpp // 可导出到你的结构体中使用 using PipelineRingBuffer::PaddingArray; // char[64] —— 手动填充用 // 示例:填充自定义数据以规避伪共享 struct MyData { double price; double volume; PaddingArray _pad; // 64 字节填充 }; ``` ## 测试 构建并运行测试: ```bash cmake -S . -B build -DPIPELINE_RINGBUFFER_BUILD_TESTS=ON cmake --build build cd build && ctest ``` | 测试名称 | 说明 | |----------|------| | `pipeline_ringbuffer.basic` | 验证数据在 4 阶段流水线中的流转(Exclusive → Exclusive → Share → Exclusive) | | `pipeline_ringbuffer.rapidjson` | 使用 rapidjson 基准测试 JSON 解析吞吐量 | | `pipeline_ringbuffer.slice_parallel_parser` | 集成测试:基于字段切片的并行 JSON 解析 | | `pipeline_ringbuffer.rapidjson_parallel_parser` | 集成测试:基于 rapidjson 的并行 JSON 解析 | | `pipeline_ringbuffer.sequence_cursor` | 验证 64 位 sequence 绕环、tail 位置和单槽位超时语义 | | `pipeline_ringbuffer.exclusive_contention` | 验证多个 ExclusiveConsumer 高竞争下不会同时持有同一槽位 | ## 许可证 MIT。详见 [LICENSE](LICENSE)。