如何用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
相关产品推荐
相关产品推荐

