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

如何用C++标准库实现带固定缓冲区的异步IO流(生产者-消费者模型)

带固定缓冲区的生产者-消费者模型实现方案

标准库原生方案说明

C++20及后续标准库没有直接提供带固定容量阻塞式的异步IO流。std::stringstream是单线程设计,不支持多线程安全的生产者-消费者模式;std::async仅用于任务调度,并非针对字节流的缓冲同步场景。

推荐现有库方案

1. Boost.Asio 现成实现

Boost.Asio的asio::streambuf完全匹配需求:

  • 可通过asio::streambuf::max_size()设置固定缓冲区上限
  • 写入时缓冲区不足会自动阻塞,读取时数据不足也会等待
  • 天然支持多线程安全,无需手动实现同步逻辑

示例代码片段:

#include <boost/asio.hpp>
#include <thread>
#include <iostream>

using namespace boost::asio;

int main() {
    io_context io;
    streambuf buf(1024); // 固定1KB缓冲区

    // 生产者线程
    std::thread producer([&]() {
        std::string data = "hello world\n";
        for (int i = 0; i < 10; ++i) {
            write(io, buf, buffer(data)); // 缓冲区满时阻塞
        }
    });

    // 消费者线程
    std::thread consumer([&]() {
        char read_buf[128];
        for (int i = 0; i < 10; ++i) {
            size_t len = read(io, buf, buffer(read_buf)); // 数据不足时阻塞
            std::cout.write(read_buf, len);
        }
    });

    producer.join();
    consumer.join();
    return 0;
}

2. C++20轻量双缓冲实现(PingPong模式)

若不想依赖Boost,用C++20的std::binary_semaphore实现双缓冲是高效简洁的选择,避免继承std::basic_streambuf的复杂逻辑:

#include <iostream>
#include <thread>
#include <semaphore>
#include <array>
#include <algorithm>

template <size_t BufferSize>
class PingPongBuffer {
private:
    std::array<char, BufferSize> buf1, buf2;
    char* active_write_buf = &buf1[0];
    char* active_read_buf = &buf2[0];
    std::binary_semaphore write_sem{1};
    std::binary_semaphore read_sem{0};
    size_t write_pos = 0;
    size_t read_pos = 0;

public:
    // 写入数据,缓冲区满则阻塞
    void write(const char* data, size_t len) {
        while (len > 0) {
            write_sem.acquire();
            size_t write_len = std::min(len, BufferSize - write_pos);
            std::copy(data, data + write_len, active_write_buf + write_pos);
            write_pos += write_len;
            data += write_len;
            len -= write_len;

            if (write_pos == BufferSize) {
                std::swap(active_write_buf, active_read_buf);
                write_pos = 0;
                read_sem.release();
            } else {
                write_sem.release();
            }
        }
    }

    // 读取数据,数据不足则阻塞
    void read(char* data, size_t len) {
        while (len > 0) {
            if (read_pos == 0) {
                read_sem.acquire();
            }
            size_t read_len = std::min(len, BufferSize - read_pos);
            std::copy(active_read_buf + read_pos, active_read_buf + read_pos + read_len, data);
            read_pos += read_len;
            data += read_len;
            len -= read_len;

            if (read_pos == BufferSize) {
                read_pos = 0;
                write_sem.release();
            }
        }
    }
};

int main() {
    PingPongBuffer<1024> buf;

    std::thread producer([&]() {
        std::string data = "hello pingpong\n";
        for (int i = 0; i < 10; ++i) {
            buf.write(data.data(), data.size());
        }
    });

    std::thread consumer([&]() {
        char read_buf[128];
        for (int i = 0; i < 10; ++i) {
            buf.read(read_buf, 12);
            std::cout.write(read_buf, 12);
        }
    });

    producer.join();
    consumer.join();
    return 0;
}

其他方案分析

  • 继承std::basic_streambuf:需手动处理线程安全和缓冲切换逻辑,极易引入bug,非必要不推荐
  • boost::lockfree::queue:无固定容量限制,且无法自动阻塞生产者/消费者,不符合需求
  • std::coroutine:仅能简化异步逻辑,仍需手动实现缓冲和同步,不属于现成解决方案,适合搭配异步框架使用

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.10 22:15:29