Asynchronous Programming with C++ 学习:Boost.Asio 的缓冲区(buffer) —— 数据到底是怎么“搬运“的
一、缓冲区是什么,为什么 I/O 操作离不开它
不管是从网络上收数据、往文件里写数据,还是给一个定时器传参数,本质上都需要一块连续的内存空间来暂存要传输的数据,这块内存就叫缓冲区(buffer)。所有的 I/O 操作说到底都是在"把数据从一块内存搬到另一块内存",缓冲区就是这个搬运过程里数据实际存放的地方。
二、两种缓冲区类型:能不能写是关键区别
Boost.Asio 把缓冲区分成两类:
| 类型 | 对应的C++类型 | 能不能写入数据 | 典型用途 |
|---|---|---|---|
| 可写缓冲区 | boost::asio::mutable_buffer | 可以(既能读也能写) | 接收数据,比如从网络读进来的内容要写进这块内存 |
| 只读缓冲区 | boost::asio::const_buffer | 不可以(只能读) | 发送数据,比如把内存里已经有的内容原样发出去 |
这两者之间有一条单向的转换规则:
mutable_buffer⟶const_buffer(合法)mutable\_buffer \longrightarrow const\_buffer \quad (\text{合法})mutable_buffer⟶const_buffer(合法)
const_buffer⟶mutable_buffer(不合法)const\_buffer \longrightarrow mutable\_buffer \quad (\text{不合法})const_buffer⟶mutable_buffer(不合法)
这个规则其实很符合直觉:一块"既能读又能写"的内存,当然可以被当成"只读"的方式去使用(反正只是不去写它而已);但反过来,一块本来就标记为"只读"的内存,是不能凭空变成"可写"的——因为它底层指向的数据可能压根就是 const 的,强行写入会破坏这个约定。
不管是哪种缓冲区,Boost.Asio 都会在内部做越界保护:操作缓冲区时不会让读写操作超出这块内存本身的大小,避免缓冲区溢出这种常见的安全隐患。
三、boost::asio::buffer() 函数:从各种数据源快速创建缓冲区
自己手写一个 mutable_buffer 或 const_buffer 对象很麻烦,Boost.Asio 提供了一个万能的辅助函数 boost::asio::buffer(),可以从好几种常见的数据类型直接构造出对应的缓冲区对象:
| 数据来源 | 示例写法 | 得到的缓冲区类型 |
|---|---|---|
| 裸指针+大小 | buffer(ptr, size) | mutable_buffer(如果指针指向非const内存) |
std::string | buffer(str) | const_buffer |
POD类型的数组/std::vector | buffer(vec) | 取决于元素是否为const |
这里提到的 POD(Plain Old Data)指的是这样一类类型:它没有用户自定义的拷贝赋值运算符、没有用户自定义的析构函数,也没有 private 或 protected 的非静态成员变量——简单说就是"结构简单、可以直接按字节拷贝"的类型,比如只包含 int、double 这类基础成员的结构体。
书里给出的最简单的例子是从一个字符数组创建可写缓冲区:
char data[1024];
mutable_buffer buffer = buffer(data, sizeof(data));
逐词说明这两行:
char data[1024];:在栈上开辟一块 1024 字节大小的内存,这是真正存放数据的地方。mutable_buffer buffer = buffer(data, sizeof(data));:调用boost::asio::buffer()函数,传入指向这块内存的指针data和它的大小sizeof(data)(也就是 102410241024 字节),得到一个mutable_buffer对象。注意这个buffer对象本身并不存储数据,它只是一个"视图"——内部记录了"数据从哪个地址开始、一共有多长",真正的数据始终存放在data这个数组里。
四、完整可运行的代码:几种常见的创建方式 + 越界保护效果
下面这份代码演示了从裸内存、std::string、POD 结构体的 vector 分别创建缓冲区,并且用 boost::asio::buffer_copy 直观展示"越界保护"到底是怎么起作用的:
#include <boost/asio.hpp>
#include <iostream>
#include <string>
#include <vector>
int main() {
// ---------- 场景一: 从裸内存指针 + 大小创建可写缓冲区 ----------
char raw_data[1024] = {};
boost::asio::mutable_buffer mbuf = boost::asio::buffer(raw_data, sizeof(raw_data));
std::cout << "raw_data缓冲区大小: " << mbuf.size() << " 字节\n";
// ---------- 场景二: 从std::string创建只读缓冲区 ----------
// std::string对象本身管理着字符数据,buffer()只是"看"着这份数据,不会拷贝
std::string message = "Hello Boost.Asio buffers!";
boost::asio::const_buffer cbuf = boost::asio::buffer(message);
std::cout << "message缓冲区大小: " << cbuf.size() << " 字节\n";
// ---------- 场景三: 从POD结构体的vector创建缓冲区 ----------
// Point满足POD条件: 没有自定义拷贝赋值/析构函数,没有private/protected成员
struct Point {
int x;
int y;
};
std::vector<Point> points = {{1, 2}, {3, 4}, {5, 6}};
boost::asio::const_buffer pbuf = boost::asio::buffer(points);
std::cout << "points缓冲区大小: " << pbuf.size() << " 字节 (等于 "
<< points.size() << " 个Point,每个 " << sizeof(Point) << " 字节)\n";
// ---------- 演示: mutable_buffer可以隐式转换成const_buffer ----------
boost::asio::const_buffer view_of_mbuf = mbuf; // 合法:可写buffer当成只读用是没问题的
std::cout << "mutable_buffer成功转换为const_buffer, 大小: "
<< view_of_mbuf.size() << " 字节\n";
// ---------- 演示: buffer_copy的越界保护效果 ----------
// 故意把26字节的message,拷贝进只有10字节大小的小缓冲区,
// buffer_copy不会越界写入,只会拷贝两者之中较小的那个长度
char small_buf[10] = {};
boost::asio::mutable_buffer small_mbuf = boost::asio::buffer(small_buf, sizeof(small_buf));
std::size_t copied = boost::asio::buffer_copy(small_mbuf, cbuf);
std::cout << "尝试把 " << cbuf.size() << " 字节拷贝进 " << sizeof(small_buf)
<< " 字节的缓冲区,实际只拷贝了 " << copied << " 字节\n";
std::cout << "小缓冲区里的内容: "
<< std::string(static_cast<const char*>(small_mbuf.data()), copied) << "\n";
return 0;
}
编译方式:
g++ -std=c++17 buffer_demo.cpp -o buffer_demo -lboost_system -lpthread
预期输出:
raw_data缓冲区大小: 1024 字节
message缓冲区大小: 26 字节
points缓冲区大小: 24 字节 (等于 3 个Point,每个 8 字节)
mutable_buffer成功转换为const_buffer, 大小: 1024 字节
尝试把 26 字节拷贝进 10 字节的缓冲区,实际只拷贝了 10 字节
小缓冲区里的内容: Hello Boost
(points 缓冲区大小 242424 字节的由来:每个 Point 结构体包含两个 int,大小是 2×4=82 \times 4 = 82×4=8 字节,一共 333 个元素,所以总大小是 3×8=243 \times 8 = 243×8=24 字节。)
buffer_copy 这一步是最能体现"越界保护"效果的地方:拷贝的字节数永远是目标缓冲区大小和源缓冲区大小两者中较小的那个,用公式表示就是:
bytes_copied=min(size(dest),size(source))bytes\_copied = \min(size(dest), size(source))bytes_copied=min(size(dest),size(source))
正因为遵循这个规则,就算源数据(26 字节)比目标缓冲区(10 字节)大得多,也绝对不会往目标缓冲区之外多写一个字节,天然避免了缓冲区溢出。
五、一个容易踩的坑:缓冲区不拥有数据,生命周期要自己管
这一点非常重要,一定要单独强调:mutable_buffer/const_buffer 对象本身只是"指针+长度"这样一个轻量的视图,它并不会复制、也不会延长底层数据的生命周期。底层内存什么时候失效、什么时候被回收,完全是程序自己的责任,Boost.Asio 库不会替你操心这件事。
下面这份代码故意演示了这个问题会带来什么后果,只用于教学演示,其中包含未定义行为,实际项目中不要这样写:
#include <boost/asio.hpp>
#include <iostream>
#include <string>
// 危险写法: 返回一个指向局部数组的buffer
boost::asio::mutable_buffer make_dangling_buffer() {
char local_data[64] = "temporary content";
// 危险点: local_data是这个函数的局部变量,函数一返回它就不存在了;
// 但buffer()构造出来的对象只是"看"着这块内存的一个视图,并没有真正拷贝内容
return boost::asio::buffer(local_data, sizeof(local_data));
}
int main() {
boost::asio::mutable_buffer dangling = make_dangling_buffer();
// 这里访问dangling背后指向的内存,其实是在访问一个已经不存在的栈变量,
// 属于未定义行为: 可能看起来内容还"正常",也可能已经变成乱码,甚至导致崩溃
std::cout << "危险访问: "
<< std::string(static_cast<char*>(dangling.data()), dangling.size())
<< std::endl;
return 0;
}
问题出在 make_dangling_buffer() 函数:local_data 是这个函数内部的局部数组,函数一返回,这块栈内存就已经失效了;但函数返回的 mutable_buffer 对象里存的仍然是指向这块内存的地址,这就是一个典型的"悬空视图"。正确的做法是保证缓冲区背后的数据(比如 std::string、std::vector,或者在外部作用域声明的数组)至少要存活到这个缓冲区被使用完毕为止,通常做法是把底层数据放在调用者那一层(比如 main 函数里)声明,而不是放在一个用完就返回的临时函数里。
六、缓冲区类型关系图
七、时序图:缓冲区在一次数据拷贝中的角色
八、小结
| 概念 | 说明 |
|---|---|
| 缓冲区(buffer) | I/O操作里用来暂存数据的连续内存区域 |
mutable_buffer | 可读可写的缓冲区,常用于接收数据 |
const_buffer | 只读缓冲区,常用于发送数据 |
boost::asio::buffer() | 从裸指针、std::string、POD数组/vector等快速构造缓冲区对象的辅助函数 |
buffer_copy() | 在两个缓冲区之间拷贝数据,只拷贝较小者的长度,天然防止越界 |
| 缓冲区的生命周期 | 缓冲区对象本身不拥有数据,底层内存必须由程序自己保证存活足够久 |
理解缓冲区机制最关键的一点是:mutable_buffer/const_buffer 只是对一块已经存在的内存的"描述",不是数据本身的容器。正因为它们这么"轻量",创建和传递缓冲区的开销很小,但代价是必须由写代码的人自己保证:只要还有异步操作在用这个缓冲区,它背后指向的那块真实内存就绝对不能被提前销毁或挪作他用。
分散-聚集操作(Scatter-Gather Operations)
一、什么是 scatter-read 和 gather-write
缓冲区(buffer)除了一个个单独使用之外,还可以把好几个缓冲区打包在一起,跟操作系统的一次读写系统调用配合使用,这种用法叫做分散-聚集操作(scatter-gather operations)。
它包含两个方向相反的操作:
- 散射读(scatter-read):数据来自同一个来源(比如同一个 socket),但读出来之后被拆分着放进好几块不连续的内存缓冲区里,而不是全部塞进一块连续的大内存。
- 聚集写(gather-write):跟散射读正好相反,数据原本分散在好几块不连续的内存缓冲区里,写出去的时候被汇总起来,一次性发送到同一个目的地(比如同一个 socket)。
二、为什么要用这种方式
如果不用分散-聚集操作,遇到"数据本来就散落在好几块内存里,但要一起发出去"这种情况,通常的做法是先把这些零散的数据拷贝、拼接成一块连续的大缓冲区,再调用一次写操作。这个"先拷贝拼接"的过程本身就是额外的开销。
用了分散-聚集操作之后,可以直接把这些零散的缓冲区打包,一次系统调用搞定,不需要先做一次额外的内存拷贝,也不需要为了凑成一块连续内存而多发起几次系统调用。这样做能明显提升效率和性能,具体带来两个好处:
- 减少了系统调用的次数(原本可能要分好几次读/写,现在一次搞定);
- 减少了数据拷贝的次数(不需要先拼接成连续内存再操作)。
这种技术不只是用在网络 I/O 上,在数据处理、机器学习、排序算法、矩阵乘法这类并行算法里也经常能看到它的身影——只要是"数据分散在多块内存里,但需要作为一个整体去处理"的场景,都可以考虑用这种思路。
用一个简单的公式来表示一次分散-聚集操作总共处理的数据量:
total_bytes=∑i=1nlen(bufi)total\_bytes = \sum_{i=1}^{n} len(buf_i)total_bytes=i=1∑nlen(bufi)
也就是说,不管数据被拆成了多少块(这里用 nnn 表示缓冲区的数量,len(bufi)len(buf_i)len(bufi) 表示第 iii 块缓冲区的长度),一次分散-聚集操作处理的总字节数,就是所有这些缓冲区长度加起来的总和。
三、怎么在 Boost.Asio 里使用
要触发分散-聚集操作,只需要把多个缓冲区打包放进一个容器里(std::vector、std::list、std::array、boost::array 都可以),然后把这个容器整体传给异步操作函数,而不是只传一个缓冲区。
用一张表把两种操作和对应的函数对应起来:
| 操作 | 数据流向 | 典型用到的容器元素类型 | 典型函数 |
|---|---|---|---|
| 散射读(scatter-read) | 从一个来源读出,分散写入多块内存 | boost::asio::mutable_buffer(可写) | socket.async_read_some(buffers, handler) |
| 聚集写(gather-write) | 从多块内存汇总,写到同一个目的地 | boost::asio::const_buffer(只读) | socket.async_write_some(buffers, handler) |
散射读用的是 mutable_buffer,因为读操作需要往这些缓冲区里写入数据,所以缓冲区必须是可修改的;聚集写用的是 const_buffer,因为写操作只需要从这些缓冲区里读出已有的数据发送出去,不需要修改它们。
四、完整代码示例:一个进程内同时演示 scatter-read 和 gather-write
下面这个例子在同一个程序里同时扮演"服务器"和"客户端"两个角色:服务器用散射读把收到的数据拆分放进两个缓冲区,客户端把两个分散的缓冲区聚集起来一次性发送出去,两边通过本地回环地址 127.0.0.1 建立 TCP 连接。
// 编译命令:
// g++ -std=c++17 scatter_gather_demo.cpp -o scatter_gather_demo -lpthread
#include <boost/asio.hpp>
#include <array>
#include <iostream>
#include <string>
#include <vector>
using boost::asio::ip::tcp;
int main() {
boost::asio::io_context io_context;
// 服务器端:在本地 55555 端口监听连接
tcp::acceptor acceptor(io_context, tcp::endpoint(tcp::v4(), 55555));
// 服务器端和客户端各自持有一个 socket,都绑定在同一个 io_context 上
tcp::socket server_socket(io_context);
tcp::socket client_socket(io_context);
// ---------- 服务器端:散射读要用到的两块内存 ----------
// 客户端总共会发 11 个字节过来:"Hello"(5字节)+ " World"(6字节)
// 服务器故意用两块大小分别匹配的缓冲区去接收,模拟"数据被拆开存放"的场景
std::array<char, 5> server_buf1{};
std::array<char, 6> server_buf2{};
// ---------- 客户端:聚集写要用到的两块内存 ----------
// 这两块内存本身并不连续(是两个独立的 std::array 变量),
// 但我们希望把它们的内容当成一个整体一次性发送出去
std::array<char, 5> client_buf1 = {'H', 'e', 'l', 'l', 'o'};
std::array<char, 6> client_buf2 = {' ', 'W', 'o', 'r', 'l', 'd'};
// 先注册服务器端的异步接受连接操作
acceptor.async_accept(server_socket,
[&](const boost::system::error_code& ec) {
if (!ec) {
std::cout << "[服务器] 接受到新连接\n";
// 把两块可写缓冲区打包进 vector,作为一个整体传给 async_read_some
// mutable_buffer 表示这块内存是可以被写入的
std::vector<boost::asio::mutable_buffer> buffers = {
boost::asio::buffer(server_buf1),
boost::asio::buffer(server_buf2)
};
// 发起散射读:操作系统会把这次收到的数据,
// 按顺序先填满 server_buf1,再继续填 server_buf2,
// 整个过程只对应一次底层的读取系统调用
server_socket.async_read_some(
buffers,
[&](const boost::system::error_code& ec2, std::size_t length) {
if (!ec2) {
std::cout << "[服务器] scatter-read 完成,"
<< "总共读取 " << length << " 字节\n";
std::cout << "[服务器] buf1 内容: "
<< std::string(server_buf1.data(), server_buf1.size())
<< "\n";
std::cout << "[服务器] buf2 内容: "
<< std::string(server_buf2.data(), server_buf2.size())
<< "\n";
} else {
std::cout << "[服务器] 读取出错: " << ec2.message() << "\n";
}
});
}
});
// 客户端发起异步连接
client_socket.async_connect(
tcp::endpoint(boost::asio::ip::make_address("127.0.0.1"), 55555),
[&](const boost::system::error_code& ec) {
if (!ec) {
std::cout << "[客户端] 连接成功,准备发起 gather-write\n";
// 把两块只读缓冲区打包进 vector,作为一个整体传给 async_write_some
// const_buffer 表示这块内存在这次操作里只会被读取,不会被修改
std::vector<boost::asio::const_buffer> buffers = {
boost::asio::buffer(client_buf1),
boost::asio::buffer(client_buf2)
};
// 发起聚集写:把 client_buf1 和 client_buf2 的内容
// 当成一份连续的数据一次性发送出去,
// 不需要我们自己先手动拼接成一块大内存再发
client_socket.async_write_some(
buffers,
[&](const boost::system::error_code& ec2, std::size_t length) {
if (!ec2) {
std::cout << "[客户端] gather-write 完成,"
<< "总共发送 " << length << " 字节\n";
} else {
std::cout << "[客户端] 写入出错: " << ec2.message() << "\n";
}
});
}
});
// 服务器和客户端的所有异步操作都注册在同一个 io_context 上,
// 一次 run() 调用就能把接受连接、散射读、聚集写这一整套流程都处理完
io_context.run();
return 0;
}
关键点说明:
- 服务器和客户端虽然写在同一个进程里,但它们通过真实的 TCP 连接(本地回环地址)通信,这套散射读/聚集写的机制跟分布在两台不同机器上的场景完全一样,只是为了方便直接运行演示才放进了一个文件。
async_write_some不保证一定会把传进去的所有数据一次性发送完——理论上它有可能只发出去一部分。这个例子里数据量很小(总共 11 字节),在本地回环连接上几乎总能一次发送完成,但如果你要保证"不管数据多大,一定把整个缓冲区序列全部发完",可以换成boost::asio::async_write(socket, buffers, handler),它会在内部自动帮你处理"只写了一部分,需要继续写"的情况。server_buf1是 5 字节,正好对应客户端发过来的"Hello";server_buf2是 6 字节,正好对应" World"。散射读会按缓冲区在容器里的顺序依次填充,先把server_buf1填满,再继续填server_buf2,这也是为什么服务器最终打印出来的两段内容能跟客户端发送时的拆分方式对应上。
五、内存布局示意:多块不连续内存如何被当成一个整体处理
用 ASCII 图直观感受一下散射读和聚集写各自发生了什么:
散射读(scatter-read):数据来源只有一个,读出来后分散存放
+-------------------------+
| socket 收到的数据流 | "Hello World"
+-------------------------+
|
操作系统一次读取
|
+--------+--------+
v v
+--------+ +---------+
| buf1 | | buf2 |
| Hello | | " World"|
+--------+ +---------+
(两块内存地址上并不相邻,只是逻辑上被当成一组一起处理)
聚集写(gather-write):数据分散在多块内存,写出去时汇总成一份
+--------+ +---------+
| buf1 | | buf2 |
| Hello | | " World"|
+--------+ +---------+
| |
+--------+--------+
|
操作系统一次写入
|
+-------------------------+
| 发送到 socket 的数据 | "Hello World"
+-------------------------+
六、时序图:完整的 gather-write 到 scatter-read 流程
七、小结
- 散射读是"一个数据来源,读出来分散存放到多块内存";聚集写是反过来,“多块分散的内存,汇总起来写到一个目的地”,两者是一对方向相反的操作。
- 用分散-聚集操作最大的好处是可以省去手动拼接内存的步骤,一次系统调用就能处理好几块不连续的缓冲区,减少了系统调用次数和数据拷贝次数,这对追求高性能的网络程序、数据处理和并行算法来说都很有意义。
- 在 Boost.Asio 里,只需要把多个缓冲区放进
std::vector、std::list、std::array或boost::array这类容器里,整体传给async_read_some(散射读用mutable_buffer)或async_write_some(聚集写用const_buffer),底层就会自动按顺序依次读写这些缓冲区。 async_write_some不保证一次性写完所有数据,如果需要"保证整个缓冲区序列全部发送完毕"这种更强的保证,应该改用boost::asio::async_write。
流缓冲区(stream buffer) 与散射操作 —— 数据长度不确定时怎么办
一、为什么需要"流缓冲区"这种东西
前面讲的 mutable_buffer/const_buffer 有一个共同特点:大小是固定的,创建的时候就得知道"这块内存有多少字节"。但现实中很多场景下,我们根本不知道即将收到的数据到底有多长——比如从网络上读一条消息,谁也没法保证对方发过来的内容正好塞得满、塞得下某个固定大小的缓冲区。
boost::asio::streambuf 就是为了解决这个问题而存在的。它是基于标准库 std::basic_streambuf(定义在 <streambuf> 头文件里)实现的一种动态缓冲区:大小可以随着实际收到的数据量自动调整,不需要提前精确计算好要开多大的内存。
二、什么是"散射-聚集"操作(scatter-gather)
正常情况下,一次 I/O 操作(比如读数据)只对应一块缓冲区。但 Boost.Asio 支持一次 I/O 操作同时操作多块缓冲区:
- 散射(scatter):读数据时,把收到的字节依次"分散"填进多个缓冲区,一个缓冲区填满了就接着填下一个。
- 聚集(gather):写数据时,把多个缓冲区里的内容"聚拢"到一起,一次性发送出去。
这样做的好处是:不需要先把所有数据拼接到一整块连续内存里,可以直接对多个分散的小缓冲区进行一次性的读写操作,减少了额外的拷贝和拼接开销。
下面这个例子会实现一个简单的 TCP 服务器,用两个streambuf配合"散射"操作,把客户端发来的数据自动分流到两个缓冲区里。为了让例子聚焦在流缓冲区和散射操作本身,这里用的是同步(阻塞)方式的 I/O,不涉及异步回调。
三、main() 函数逐行解析:搭建一个最简单的 TCP 服务器
boost::asio::io_context io_context;
tcp::acceptor acceptor(io_context, tcp::endpoint(tcp::v4(), port));
std::cout << "Server is running on port " << port << "...\n";
while (true) {
tcp::socket socket(io_context);
acceptor.accept(socket);
std::cout << "Client connected...\n";
handle_client(socket);
std::cout << "Client disconnected...\n";
}
tcp::acceptor acceptor(io_context, tcp::endpoint(tcp::v4(), port));:创建一个"接待员"对象,负责监听指定端口(这里是1234),等着接收客户端的连接请求。tcp::v4()表示使用 IPv4 协议。while (true) { ... }:一个死循环,服务器会持续不断地接受新客户端连接,一个处理完了就继续等下一个。tcp::socket socket(io_context);:为即将到来的这个客户端连接准备一个空的 socket 对象。acceptor.accept(socket);:阻塞等待,直到有客户端真正连接进来,把这个连接绑定到socket上。handle_client(socket);:调用处理函数,专门负责跟这个客户端交互(读取它发来的数据)。这个函数处理完、客户端断开连接之后,循环会继续,回去等待下一个客户端。
四、handle_client() 函数逐行解析:核心的流缓冲区+散射操作
void handle_client(tcp::socket& socket) {
const size_t size_buffer = 5;
boost::asio::streambuf buf1, buf2;
std::array<boost::asio::mutable_buffer, 2> buffers = {
buf1.prepare(size_buffer),
buf2.prepare(size_buffer)
};
boost::system::error_code ec;
size_t bytes_recv = socket.read_some(buffers, ec);
if (ec) {
std::cerr << "Error on receive: " << ec.message() << '\n';
return;
}
std::cout << "Received " << bytes_recv << " bytes\n";
buf1.commit(5);
buf2.commit(5);
std::istream is1(&buf1);
std::istream is2(&buf2);
std::string data1, data2;
is1 >> data1;
is2 >> data2;
std::cout << "Buffer 1: " << data1 << std::endl;
std::cout << "Buffer 2: " << data2 << std::endl;
}
逐段拆解:
boost::asio::streambuf buf1, buf2;:创建两个流缓冲区,此时它们都是空的,还没有分配任何具体大小的存储空间。buf1.prepare(size_buffer):prepare(n)的作用是"预留"出n个字节(这里是 5 字节)的可写空间,并返回一个指向这块空间的mutable_buffer。注意这一步只是预留空间,还没有真正确认里面写进了多少有效数据——这是流缓冲区特有的"两阶段"写入方式:先prepare()要一块空间,写完数据后再调用commit()才算数。std::array<boost::asio::mutable_buffer, 2> buffers = {...};:把两个prepare()返回的可写缓冲区放进一个数组里,凑成一组"缓冲区序列",接下来把这一整组传给读操作,就能实现"散射"效果。socket.read_some(buffers, ec);:同步地从 socket 读取数据。因为传入的是包含两个缓冲区的数组,Boost.Asio 会自动按顺序先填满第一个缓冲区(buf1预留的 5 字节),如果还有剩余数据,接着填第二个缓冲区(buf2预留的 5 字节)——这就是"散射"操作的具体体现。返回值bytes_recv是这次操作实际读到的总字节数;如果出错,错误信息会写进ec里,而不会抛出异常。- 错误处理:如果
ec表示有错误,打印错误信息并直接返回,不再继续往下处理。 buf1.commit(5); buf2.commit(5);:commit(n)的作用是"确认"——告诉流缓冲区"刚才prepare()出来的那块空间里,前n个字节现在是真正有效的数据了,可以被当成’可读内容’对待"。只有commit()之后,这些数据才能被后续的读取操作(比如下面的istream)看到。std::istream is1(&buf1); std::istream is2(&buf2);:把流缓冲区包装成标准的输入流对象,这样就能用熟悉的>>运算符从里面提取内容,跟读std::cin的写法几乎一样。is1 >> data1; is2 >> data2;:用流提取运算符>>分别从两个流缓冲区里读出一个"单词"(>>默认会跳过开头的空白字符,读到下一个空白字符或者数据末尾为止)。
五、完整可运行的代码
把上面两部分拼在一起,就是一份完整、可以直接编译运行的 TCP 服务器程序:
#include <array>
#include <iostream>
#include <string>
#include <boost/asio.hpp>
#include <boost/asio/streambuf.hpp>
using boost::asio::ip::tcp;
constexpr int port = 1234;
// 提前声明,因为main()里要用到它,但函数体写在main()后面
void handle_client(tcp::socket& socket);
int main() {
try {
boost::asio::io_context io_context;
tcp::acceptor acceptor(io_context, tcp::endpoint(tcp::v4(), port));
std::cout << "Server is running on port " << port << "...\n";
while (true) {
tcp::socket socket(io_context);
acceptor.accept(socket); // 阻塞,直到有客户端连进来
std::cout << "Client connected...\n";
handle_client(socket);
std::cout << "Client disconnected...\n";
}
} catch (std::exception& e) {
std::cerr << "Exception: " << e.what() << '\n';
}
return 0;
}
void handle_client(tcp::socket& socket) {
const size_t size_buffer = 5;
boost::asio::streambuf buf1, buf2;
// 各自预留5字节的可写空间,凑成一组缓冲区序列,用于散射读取
std::array<boost::asio::mutable_buffer, 2> buffers = {
buf1.prepare(size_buffer),
buf2.prepare(size_buffer)
};
boost::system::error_code ec;
// 一次read_some调用,数据会按顺序自动分散填进buf1、buf2
size_t bytes_recv = socket.read_some(buffers, ec);
if (ec) {
std::cerr << "Error on receive: " << ec.message() << '\n';
return;
}
std::cout << "Received " << bytes_recv << " bytes\n";
// 确认写入的字节数正式生效,变成可读内容
buf1.commit(5);
buf2.commit(5);
std::istream is1(&buf1);
std::istream is2(&buf2);
std::string data1, data2;
is1 >> data1;
is2 >> data2;
std::cout << "Buffer 1: " << data1 << std::endl;
std::cout << "Buffer 2: " << data2 << std::endl;
}
编译方式:
g++ -std=c++17 streambuf_server.cpp -o streambuf_server -lboost_system -lpthread
六、怎么运行这个例子
需要开两个终端:
- 第一个终端运行编译好的服务器程序:
./streambuf_server - 第二个终端用
telnet连接上去发消息:telnet 127.0.0.1 1234,连上之后输入一段内容,比如Hello World,回车发送。
服务器终端会打印类似下面的内容:
Server is running on port 1234...
Client connected...
Received 10 bytes
Buffer 1: Hello
Buffer 2: Worl
Client disconnected...
七、为什么最终只处理了 10 字节,中间的空格又去哪了
"Hello World" 这段文字一共有 11 个字符(H e l l o (空格) W o r l d),但由于两个缓冲区各自只预留了 5 字节,加起来总共能装下的数据量是:
5+5=10 字节5 + 5 = 10 \text{ 字节}5+5=10 字节
所以 read_some() 这一次只会读取前 10 个字节,剩下的最后一个字符 'd' 根本没有被这次调用读到。这 10 个字节按顺序被"散射"分配成:
输入字节(前10个,用下划线表示空格): H e l l o _ W o r l
|_________| |_________|
buf1(5字节) buf2(5字节)
"Hello" " Worl"
buf1 拿到的是 "Hello",buf2 拿到的是 " Worl"(开头带一个空格)。之后用 is2 >> data2 提取时,>> 运算符会自动跳过开头的空白字符,所以最终 data2 变成了 "Worl"——这也是为什么最后打印出来的 Buffer 2 缺了开头那个空格、也缺了最后那个没被读到的 'd'。
八、prepare() 和 commit() 之间的关系
| 方法 | 作用 |
|---|---|
prepare(n) | 预留n字节的可写空间,返回对应的mutable_buffer,此时数据还不算"生效" |
commit(n) | 确认前n字节已经写好,正式把它们纳入可以被读取的范围 |
istream提取 | 只能读到已经commit()过的内容,prepare()但未commit()的部分读不到 |
可以把这个过程理解成两步走:先申请一块空地(prepare),把东西放进去,再去登记造册(commit)确认这块地上的东西可以对外公开使用。只申请不登记,东西虽然物理上已经写进内存了,但在"可读取"的意义上依然是不存在的。
九、时序图:一次完整的读取过程
十、流缓冲区适合什么场景
| 场景 | 适合用固定大小缓冲区 | 适合用流缓冲区(streambuf) |
|---|---|---|
| 提前知道数据长度 | 适合 | 也可以,但没必要 |
| 数据长度不固定、不可预知 | 不适合,容易装不下或浪费空间 | 非常适合,会按需自动扩展 |
| 想减少内存拷贝开销 | 需要自己精心设计 | prepare/commit机制本身就是为此设计的 |
流缓冲区和固定大小的缓冲区并不是互斥的关系,两者完全可以配合使用——比如本例就是把两个流缓冲区放进同一个数组里,一起参与散射读取操作。
十一、小结
| 概念 | 作用 |
|---|---|
boost::asio::streambuf | 大小可以动态调整的缓冲区,适合处理长度不确定的数据 |
| 散射(scatter) | 一次读操作把数据依次分流进多个缓冲区 |
| 聚集(gather) | 一次写操作把多个缓冲区的内容合并发送 |
prepare(n) | 预留n字节可写空间 |
commit(n) | 确认前n字节正式生效,变成可读内容 |
std::istream(&streambuf) | 把流缓冲区包装成标准输入流,方便用>>提取内容 |
这个例子最值得记住的一点是:流缓冲区把"数据长度不确定"这个麻烦事,转换成了"先预留、用多少就确认多少"这样一套简单的两阶段流程,配合散射操作,就可以很自然地把一段连续的输入数据,按需分流进多个不同用途的缓冲区里,而不需要提前精确计算好每一块该开多大。
Boost.Asio 信号处理详解
一、这一节要解决什么问题
我们写的服务器或者后台程序,通常不是自己想退出就退出的,而是要等操作系统"通知"它退出。比如你在终端按下 Ctrl+C,操作系统会给这个进程发一个叫 SIGINT 的信号;如果用 kill 命令或者容器编排系统(比如 Kubernetes、Docker)想要停掉一个进程,通常会先发 SIGTERM 信号,礼貌地让程序自己收尾,如果程序不理睬,过一段时间才会强制杀掉(SIGKILL,这个信号是不能被捕获的)。
所谓"优雅关闭"(graceful shutdown),说的就是:程序收到这种"请你退出"的信号后,不是立刻被操作系统一刀切死,而是有机会先把手头的事情收拾好——比如把还没写完的日志刷到磁盘、把数据库连接关掉、把正在处理的网络请求处理完——然后再正常退出。
Boost.Asio 提供了一个专门的类,叫 boost::asio::signal_set,就是用来"监听"这些信号的。它的用法和 Asio 里其他异步操作(比如异步读、异步写、定时器)是一套逻辑:发起一个"异步等待",指定一个"完成后要执行的回调函数",然后信号真的发生时,Asio 的事件循环会自动调用这个回调。
二、核心概念一步一步拆开讲
2.1 信号(Signal)是什么
信号是操作系统内核发给进程的一种"软件中断",用来通知进程发生了某个事件。常见的几个:
| 信号名 | 常见触发方式 | 默认行为 | 能否被程序捕获 |
|---|---|---|---|
SIGINT | 终端按下 Ctrl+C | 终止进程 | 能 |
SIGTERM | kill <pid>、容器停止时发出 | 终止进程 | 能 |
SIGKILL | kill -9 <pid> | 强制终止 | 不能捕获 |
SIGHUP | 终端关闭、常用于"重新加载配置" | 终止进程 | 能 |
我们这节课关心的就是前两个:SIGINT 和 SIGTERM,因为它们是"可以被程序捕获并做善后处理"的信号。
2.2 signal_set 这个类到底在做什么
boost::asio::signal_set 本质上是一个"信号收件箱"。构造的时候,你告诉它:“我要关注 io_context(asio 的事件循环),并且我要关注 SIGINT 和 SIGTERM 这两种信号”。
然后你调用它的 async_wait() 方法,意思是:“现在开始蹲着等,只要上面两个信号中的任意一个发生了,就调用我给你的这个回调函数。” 这个过程是异步的——也就是说,调用 async_wait() 这一行代码本身立刻就返回了,不会卡住程序,真正等待信号是在后台(io_context 的事件循环)里进行的。
这跟异步读写 socket 是完全一样的套路,可以类比成:
socket.async_read_some(buffer, handler)—— “数据到了就告诉我”signals.async_wait(handler)—— “信号来了就告诉我”
两者都遵循 Asio 统一的"发起异步操作 + 注册完成回调"模式。
2.3 为什么一定要有 io_context.run()
这里有一个初学者容易犯迷糊的地方:光调用 async_wait() 是不会真的去等信号的。async_wait() 只是把这个等待任务"注册"到 io_context 里,真正让这个事件循环"转起来"、去检测信号是否发生、并且在发生时调用回调函数的,是 io_context.run() 这一行。
可以理解成:
async_wait()= 往任务清单上写一条"提醒我"io_context.run()= 真正开始不停地检查任务清单,一旦有任务该被触发了,就去执行对应的回调
io_context.run()会一直阻塞(占用当前线程),直到io_context里所有的异步任务都完成,或者有人手动调用了io_context.stop()。
三、完整可运行的 C++ 代码(逐行注释版)
下面这份代码可以直接编译运行,试一下 Ctrl+C 就能看到效果。编译时需要链接 Boost.Asio 相关库(如果是纯 header-only 用法,通常只需要 -lpthread,某些平台还需要 -lboost_system,取决于你的 Boost 版本和是否定义了 BOOST_ASIO_STANDALONE)。
// signal_handling_demo.cpp
// 演示 Boost.Asio 如何优雅地捕获 SIGINT / SIGTERM 信号并退出程序
#include <boost/asio.hpp> // Boost.Asio 核心头文件,包含 io_context、signal_set 等
#include <iostream> // 用于标准输入输出(std::cout、std::cerr)
int main() {
try {
// -----------------------------------------------------------------
// 第一步:创建 io_context 对象。
// io_context 是 Asio 的"事件循环核心",所有异步操作(定时器、网络、
// 信号等)最终都是靠它来调度和执行回调函数的。
// -----------------------------------------------------------------
boost::asio::io_context io_context;
// -----------------------------------------------------------------
// 第二步:创建 signal_set 对象,绑定到 io_context 上,
// 并且告诉它我们要关注哪些信号:这里是 SIGINT(Ctrl+C)
// 和 SIGTERM(比如被 kill 命令终止时发出的信号)。
// 这个构造函数的重载支持一次传入多个信号编号。
// -----------------------------------------------------------------
boost::asio::signal_set signals(io_context, SIGINT, SIGTERM);
// -----------------------------------------------------------------
// 第三步:定义一个 lambda 表达式作为"信号发生后要执行的回调函数"。
// 这个回调函数的函数签名是 Asio 规定好的固定格式:
// void(const boost::system::error_code& ec, int signal_number)
// - ec:如果等待过程本身出错(而不是信号触发),这里会有错误码;
// 如果一切正常、确实收到了信号,ec 是"无错误"状态。
// - signal 参数:具体是哪个信号触发了这次回调(比如 2 代表 SIGINT)。
//
// 这里用 [&] 按引用捕获外部变量(主要是为了能在回调里访问
// io_context,从而调用 io_context.stop())。
// -----------------------------------------------------------------
auto handle_signal = [&](
const boost::system::error_code& ec,
int signal) {
// 判断这次唤醒是不是因为真的收到了信号(而不是异常错误)
if (!ec) {
std::cout << "收到信号: " << signal << std::endl;
// ---------------------------------------------------------
// 这里就是"优雅关闭"要做的事情:
// 在真实项目中,这一行往下可以写:
// - 停止接受新的网络连接
// - 等待正在处理的请求处理完
// - 关闭数据库连接、刷新日志缓冲区
// - 释放各种资源
// 本例只是个演示,所以直接跳到收尾。
// ---------------------------------------------------------
// 调用 io_context.stop() 会让正在运行的 io_context.run()
// 尽快返回(它不会立刻掐断当前正在执行的回调,但会让
// run() 停止继续调度新的异步任务)。
io_context.stop();
}
};
// -----------------------------------------------------------------
// 第四步:发起异步等待。
// 注意:这一行代码本身是"非阻塞"的,调用完立刻就往下执行了,
// 真正的等待动作发生在后面 io_context.run() 里。
// -----------------------------------------------------------------
signals.async_wait(handle_signal);
std::cout << "程序正在运行中,按下 Ctrl+C 可以停止程序...\n";
// -----------------------------------------------------------------
// 第五步:启动事件循环。
// 这一行会阻塞当前线程,直到:
// (a) io_context 里所有异步任务都执行完了;或者
// (b) 有人调用了 io_context.stop()(就是我们在回调里做的事)。
// 这里因为只注册了一个异步等待任务(信号等待),所以在
// 收到信号、调用 stop() 之前,run() 会一直"卡"在这里,
// 但这不是傻等,而是高效地睡眠等待操作系统事件通知。
// -----------------------------------------------------------------
io_context.run();
std::cout << "程序已经干净地退出。\n";
} catch (std::exception& e) {
// 捕获运行过程中可能抛出的标准异常(例如 io_context 构造失败等)
std::cerr << "异常: " << e.what() << '\n';
}
return 0;
}
编译示例(Linux,g++):
g++ -std=c++17 signal_handling_demo.cpp -o signal_demo -lboost_system -lpthread
./signal_demo
# 然后按 Ctrl+C,你会看到打印出信号编号,然后程序干净退出
3.1 代码里几个关键点单独说明
关键点 1:async_wait 的回调签名是固定的
Asio 里几乎所有异步操作的回调,第一个参数永远是 const boost::system::error_code&,代表这次操作有没有出错。对于 signal_set 来说,第二个参数额外多了一个 int signal,专门告诉你"具体是哪个信号触发的"——因为一个 signal_set 可以同时监听多种信号,回调里就需要知道到底是哪一种发生了,才能针对性处理(比如 SIGHUP 可能是要重新加载配置,而 SIGTERM 是要退出)。
关键点 2:为什么用 [&] 而不是 [=] 捕获
因为回调函数里需要调用 io_context.stop(),io_context 是 main() 函数里的局部变量。如果按值捕获([=]),lambda 内部拿到的会是这个对象的一个拷贝(而且 io_context 本身通常也不允许拷贝),这样调用 stop() 根本停不了外面真正在跑的那个事件循环。所以必须用引用捕获 [&],确保操作的是同一个对象。
关键点 3:io_context.stop() 之后到底发生了什么
调用 stop() 不会立刻打断当前正在执行的回调函数本身(也就是说 handle_signal 函数体里 stop() 后面如果还有代码,还是会执行完),但它会让 io_context.run() 尽快返回,不再去处理其他排队中的异步任务。这就是为什么示例里 io_context.run() 后面那句"程序已经干净地退出"能够被打印出来——run() 真的返回了。
关键点 4:同步等待的写法
书里提到还可以用 signals.wait()(不带 async_ 前缀)来同步等待信号,也就是说这一行代码会直接卡住,一直等到信号发生才往下走,不需要配合 io_context.run()。这种写法比较简单粗暴,适合"只想单纯等一个信号,其他啥都不干"的场景,不需要事件循环的其他功能。
关键点 5:多线程场景下的限制
如果你的程序是多线程的(比如好几个线程一起调用 io_context.run() 来并发处理任务,这是 Asio 常见的"线程池"用法),那么处理信号的回调函数,必须在跟 io_context 对象绑定的同一个线程里执行,通常约定就是主线程。这是操作系统层面的限制:信号的分发本身在多线程程序里就比较微妙,POSIX 标准里信号可能被发到进程内任意一个没有屏蔽该信号的线程,Boost.Asio 为了让行为可预测,内部做了处理,但代价就是要求你的信号等待逻辑跑在固定的线程上。
四、执行流程时序图(Mermaid)
下面这张时序图展示了从程序启动、注册信号监听,到用户按下 Ctrl+C、程序优雅退出的完整过程:
五、整体结构 ASCII 图
再用一张简单的文本图,展示各个对象之间"谁包含谁、谁调用谁"的关系:
main()
|
|-- 拥有 --> io_context (事件循环核心)
| ^
| | 绑定在同一个 io_context 上
| |
|-- 拥有 --> signal_set (监听 SIGINT, SIGTERM)
| |
| |-- async_wait(handle_signal)
| | |
| | v
| | handle_signal 回调
| | |
| | |-- 打印信号编号
| | |-- 执行清理逻辑
| | `-- io_context.stop()
| |
`-- 调用 --> io_context.run() <-- 真正驱动事件循环转起来的那一行
|
`-- 阻塞直到 stop() 被调用或所有任务完成
六、和其他异步操作的对比小结
为了帮助理解"信号等待"在 Asio 体系里的位置,可以对比一下它和定时器、socket 读写的相似性:
| 异步操作 | 发起方法 | 回调参数 | 驱动力 |
|---|---|---|---|
| 信号等待 | signals.async_wait(handler) | (error_code, int signal) | io_context.run() |
| 定时器 | timer.async_wait(handler) | (error_code) | io_context.run() |
| socket 读 | socket.async_read_some(buf, handler) | (error_code, size_t bytes) | io_context.run() |
可以看出,Asio 里所有异步操作都遵循同一套设计哲学:“注册一个操作 + 一个回调”,然后统一交给 io_context.run() 去驱动。一旦掌握了这个套路,学会一种异步操作(比如信号等待),基本上就等于学会了怎么用其他所有异步操作。
七、小结
- 操作系统通过信号(如
SIGINT、SIGTERM)通知进程"该退出了",程序如果什么都不做,会被系统强行终止,来不及做收尾工作。 boost::asio::signal_set让我们可以在 Asio 的异步框架内"捕获"这些信号,而不是让它们直接杀死进程。- 使用方式是标准的 Asio 套路:构造
signal_set绑定io_context和关心的信号编号 → 调用async_wait()注册回调 → 调用io_context.run()驱动事件循环,真正开始监听。 - 回调函数里可以做清理工作,最后通常调用
io_context.stop()让事件循环退出,从而让run()返回,程序自然结束。 - 也存在同步版本的
signals.wait(),直接阻塞等待,不依赖io_context.run()。 - 多线程程序中,信号相关的回调必须运行在跟
io_context绑定的同一个线程(通常是主线程)里。
Boost.Asio 中的 strand:不用锁也能让异步任务按顺序执行
1. 先搞清楚:为什么需要 strand
在写异步程序的时候,我们经常会遇到这样的问题:多个异步操作的完成回调(completion handler)可能会"同时"被触发,如果这些回调都要去动同一份数据(比如同一个文件、同一个容器),那就存在竞态条件(race condition),传统做法是用互斥锁(mutex)把这段代码锁起来。
但加锁这件事本身是有代价的:容易死锁、容易忘记解锁、多线程调试起来很痛苦。Boost.Asio 提供了另一种思路——strand。
strand 翻译过来是"链、绳",在这里可以理解成一根"队列绳子":凡是通过它提交(post)的任务,不管来自多少个不同的线程,都会被排成一队,一个接一个地严格串行执行,绝不会有两个任务同时跑。这样一来,只要你把所有会碰同一份数据的操作都通过同一个 strand 提交,就完全不需要互斥锁了。
strand 分两种:
- 隐式(implicit)strand:不需要你显式创建任何东西,代码结构本身就保证了顺序执行。
- 显式(explicit)strand:你主动创建一个
strand对象,把需要串行化的任务都通过它 post 出去。
下面分别来看。
2. 隐式 strand
2.1 情况一:只用一个线程跑 io_context
如果你的程序只调用一次 io_context::run(),并且只用一个线程去跑它,那么所有的事件处理函数天然就是排队执行的——因为压根就只有一个线程在处理,不可能出现两个处理函数同时执行的情况。这就是最简单的隐式 strand。
2.2 情况二:链式异步调用(定时器自我递归的例子)
第二种隐式 strand 出现在"一个异步操作的回调里,又发起了下一个异步操作"这种链式结构中。下面这个例子是一个会不停重启自己的定时器:每次到期后打印一条消息,然后重新设置 1 秒后再次到期,如此循环。
#include <boost/asio.hpp>
#include <chrono>
#include <iostream>
using namespace std::chrono_literals;
// 定时器到期后被调用的处理函数
// 参数:timer——引用外部创建好的定时器对象;count——当前是第几次触发,用来打印计数
void handle_timer_expiry(boost::asio::steady_timer& timer, int count) {
std::cout << "Timer expired. Count: " << count << std::endl;
// 把定时器的到期时间往后推1秒(从现在这一刻开始计时,而不是从上一次到期时间开始)
timer.expires_after(1s);
// 重新发起一次异步等待。
// lambda 按引用捕获 timer(因为定时器对象本身要保持存活、可持续复用),
// 按值捕获 count(每次调用都要用当时的计数值,不希望被后续修改影响)
timer.async_wait([&timer, count](const boost::system::error_code& ec) {
if (!ec) {
// 没有出错,说明确实是等待到期触发的,而不是被取消
// 递归调用自己,count+1,形成"链式异步调用"
handle_timer_expiry(timer, count + 1);
} else {
std::cerr << "Error: " << ec.message() << std::endl;
}
});
}
int main() {
boost::asio::io_context io_context;
// 创建一个稳定时钟定时器,初始1秒后到期
boost::asio::steady_timer timer(io_context, 1s);
int count = 0;
// 第一次发起异步等待,触发后进入 handle_timer_expiry 开始"自我循环"
timer.async_wait([&](const boost::system::error_code& ec) {
if (!ec) {
handle_timer_expiry(timer, count);
} else {
std::cerr << "Error: " << ec.message() << std::endl;
}
});
// 启动事件循环。注意:这里会一直运行下去(定时器不停自我重启),
// 如果只是想看效果,可以把这一行换成 io_context.run_for(5s); 跑5秒就退出
io_context.run();
return 0;
}
代码逐段讲解:
handle_timer_expiry是核心:它先打印当前计数,然后调用timer.expires_after(1s)把定时器的到期时间重新设置为"从现在起1秒后"。- 紧接着又调用了一次
timer.async_wait(...),这次传入的 lambda 里,一旦触发成功(!ec为真),就再次调用handle_timer_expiry自己,并把计数加一——这就是所谓"链式异步操作":一个异步操作的回调里发起了下一个异步操作,层层相扣。 - 因为整个程序只用了一个线程去调用
io_context.run(),所以不管这个链条循环多少次,每一次handle_timer_expiry的执行都不会和另一次重叠,这就是"隐式 strand"的效果——完全不需要锁,顺序性由代码结构自然保证。
运行效果是每隔1秒打印一行:
Timer expired. Count: 0
Timer expired. Count: 1
Timer expired. Count: 2
...
3. 显式 strand:多线程写同一个日志文件
当隐式 strand 满足不了需求的时候(比如你就是要用多个线程并行处理任务,但其中某一部分操作必须串行化),就要用到显式 strand了,也就是 boost::asio::strand 对象。
下面这个例子会创建一个 Logger(日志记录器)类,让 4 个线程各自写 5 条消息到同一个日志文件里,完全不使用互斥锁,靠 strand 来保证写文件操作不会互相打架。
3.1 Logger 类的实现
#include <boost/asio.hpp>
#include <chrono>
#include <fstream>
#include <iostream>
#include <memory>
#include <sstream>
#include <string>
#include <thread>
#include <vector>
using namespace std::chrono_literals;
class Logger {
public:
// 构造函数:传入 io_context(用来创建strand)和日志文件名
Logger(boost::asio::io_context& io_context, const std::string& filename)
: strand_(boost::asio::make_strand(io_context)), // 基于 io_context 的执行器创建一个 strand
file_(filename, std::ios::out | std::ios::app) // 以"追加"模式打开文件,不存在则自动创建
{
if (!file_.is_open()) {
// 如果文件打开失败(比如路径不存在、没权限),直接抛异常终止构造
throw std::runtime_error("Failed to open log file");
}
}
// 公共接口:任何线程都可以调用这个函数来写一条日志
void log(const std::string message) {
// 关键点:这里不是直接调用 do_log 去写文件,
// 而是把"写文件"这个操作打包成一个任务,post 到 strand_ 上去排队。
// 不管多少个线程同时调用 log(),它们提交的任务都会被 strand 排成一条队列,
// 依次取出执行,绝不会有两个 do_log 同时跑,因此不需要互斥锁。
boost::asio::post(strand_, [this, message]() {
do_log(message);
});
}
private:
// 真正执行写文件操作的私有函数,只会被 strand 串行调用
void do_log(const std::string& message) {
file_ << message << std::endl;
}
boost::asio::strand<boost::asio::io_context::executor_type> strand_; // 串行化队列
std::ofstream file_; // 输出文件流
};
这里最重要的两个细节:
log()里为什么按值捕获message?
因为log()是从工作线程里被调用的,它传进来的message是个临时字符串。如果 lambda 按引用捕获它,那么等log()函数返回、这个临时字符串被销毁之后,strand 队列里排队等待执行的任务引用的就是一块已经失效的内存——等真正执行do_log的时候,读到的就是垃圾数据(乱码、内容不完整)。按值捕获会复制一份message,让这份数据的生命周期跟着 lambda 走,直到任务真正执行完才销毁,这样就安全了。- 为什么
post到 strand 就不需要锁?
strand内部维护了一条任务队列,不管这些任务是从哪个线程 post 进来的,strand 保证:同一时刻,最多只有一个任务在执行,并且按照它们被 post 进来的顺序依次执行。这正好就是互斥锁想达到的效果(“同一时刻只有一个人能进临界区”),但 strand 是通过"任务排队"而不是"线程阻塞"来实现的,所以更轻量,也不会有死锁的风险。
3.2 工作线程函数
const unsigned num_messages_per_thread = 5;
// 每个线程要执行的任务:循环5次,每次拼一条消息交给 logger 记录
// 参数:logger——共享的日志器实例;id——线程编号,用于区分消息来源
void worker(std::shared_ptr<Logger> logger, int id) {
for (unsigned i = 0; i < num_messages_per_thread; ++i) {
std::ostringstream oss;
oss << "Thread " << id << " logging message " << i;
logger->log(oss.str()); // 提交日志任务(内部会post到strand)
std::this_thread::sleep_for(100ms); // 睡100毫秒,故意制造线程交替执行的效果,便于观察串行化结果
}
}
worker 函数很直白:每个线程各自跑一份,循环5次,每次生成一条带编号的消息字符串,调用 logger->log() 提交出去,然后睡 100 毫秒。这个睡眠是特意加的,目的是让 4 个线程的写日志动作在时间上交替发生,这样才能清楚地看出 strand 到底是怎么把它们排队的。
3.3 main 函数
const std::string log_filename = "log.txt";
const unsigned num_threads = 4;
int main() {
try {
boost::asio::io_context io_context;
// work_guard 的作用:告诉 io_context "还有活要干,先别退出"。
// 因为 strand 里的任务是异步 post 进去的,如果没有 work_guard,
// io_context.run() 可能在工作线程还没来得及 post 任务时就以为没事干了,提前退出。
auto work_guard = boost::asio::make_work_guard(io_context);
// 用 shared_ptr 创建唯一一个 Logger 实例,所有线程共享同一个对象、同一个文件
auto logger = std::make_shared<Logger>(io_context, log_filename);
std::cout << "Logging " << num_messages_per_thread
<< " messages from " << num_threads << " threads\n";
// 创建工作线程池:4个线程各自跑 worker() 函数
std::vector<std::jthread> threads;
for (unsigned i = 0; i < num_threads; ++i) {
threads.emplace_back(worker, logger, i);
}
// 再额外起一个线程专门跑 io_context 的事件循环,
// 让它跑2秒——因为我们知道所有消息肯定能在2秒内处理完
threads.emplace_back([&]() {
io_context.run_for(2s);
});
// std::jthread 会在析构时自动 join,所以这里 vector<jthread> 销毁时
// 会自动等待所有线程执行完毕(离开这个作用域时发生)
} catch (std::exception& e) {
std::cerr << "Exception: " << e.what() << '\n';
}
std::cout << "Done!" << std::endl;
return 0;
}
3.4 完整可编译运行的代码
把以上三部分拼在一起,就是一份完整可以直接编译运行的程序(需要链接 Boost.Asio,C++20 支持 std::jthread):
#include <boost/asio.hpp>
#include <chrono>
#include <fstream>
#include <iostream>
#include <memory>
#include <sstream>
#include <string>
#include <thread>
#include <vector>
using namespace std::chrono_literals;
// ------------------ Logger类:用strand串行化写文件操作 ------------------
class Logger {
public:
Logger(boost::asio::io_context& io_context, const std::string& filename)
: strand_(boost::asio::make_strand(io_context)),
file_(filename, std::ios::out | std::ios::app)
{
if (!file_.is_open()) {
throw std::runtime_error("Failed to open log file");
}
}
void log(const std::string message) {
boost::asio::post(strand_, [this, message]() {
do_log(message);
});
}
private:
void do_log(const std::string& message) {
file_ << message << std::endl;
}
boost::asio::strand<boost::asio::io_context::executor_type> strand_;
std::ofstream file_;
};
// ------------------ 工作线程函数 ------------------
const unsigned num_messages_per_thread = 5;
void worker(std::shared_ptr<Logger> logger, int id) {
for (unsigned i = 0; i < num_messages_per_thread; ++i) {
std::ostringstream oss;
oss << "Thread " << id << " logging message " << i;
logger->log(oss.str());
std::this_thread::sleep_for(100ms);
}
}
// ------------------ main函数 ------------------
const std::string log_filename = "log.txt";
const unsigned num_threads = 4;
int main() {
try {
boost::asio::io_context io_context;
auto work_guard = boost::asio::make_work_guard(io_context);
auto logger = std::make_shared<Logger>(io_context, log_filename);
std::cout << "Logging " << num_messages_per_thread
<< " messages from " << num_threads << " threads\n";
std::vector<std::jthread> threads;
for (unsigned i = 0; i < num_threads; ++i) {
threads.emplace_back(worker, logger, i);
}
threads.emplace_back([&]() {
io_context.run_for(2s);
});
} catch (std::exception& e) {
std::cerr << "Exception: " << e.what() << '\n';
}
std::cout << "Done!" << std::endl;
return 0;
}
编译命令示例(Linux,已安装 Boost):
g++ -std=c++20 -O2 -pthread strand_logger.cpp -o strand_logger -lboost_system
./strand_logger
4. 运行结果对比:三种不同写法,结果天差地别
同样是这段程序,如果对细节做一点改动,写出来的 log.txt 内容会完全不一样。下面用表格对比三种情况:
| 情况 | 代码改动 | log.txt 的内容顺序 | 原因 |
|---|---|---|---|
| 原始版本 | 保留 work_guard,每次写完睡100ms | 按"消息序号"排序:先是4个线程各自的第0条,再是各自的第1条…… | 4个线程写消息的节奏因为sleep而错开,谁先post谁先被strand执行 |
| 去掉 work_guard | 其余不变 | 只有4个线程各自的第0条消息 | io_context 在4个线程刚各post完第一条消息、还没来得及跑第二轮时,就以为没活干了,提前退出 |
| 去掉 work_guard 且去掉 sleep_for | worker函数里不再睡眠 | 按"线程编号"排序:线程0的5条全部先写完,然后是线程1的5条…… | 没有睡眠时,每个线程会几乎瞬间把自己的5条消息全部post出去,谁先抢到CPU时间片,谁的5条就会连续排在strand队列前面 |
从这个对比能看出两件事:
- strand 只保证"串行、不交叉",不保证"公平地轮流执行"——具体谁先谁后,取决于任务被 post 进队列的实际时间点。
work_guard很关键——只要io_context觉得"手头队列空了、没有更多工作要来了",它就会认为可以退出run()/run_for(),哪怕别的线程随后还会 post 新任务进来。work_guard相当于告诉它"先别急着走,后面还有事"。
5. 一个容易踩的坑:lambda 按值捕获 vs 按引用捕获
在 log() 函数里,我们是这样 post 任务的:
boost::asio::post(strand_, [this, message]() { do_log(message); });
这里 message 是按值捕获的。如果图省事改成按引用捕获:
boost::asio::post(strand_, [&]() { do_log(message); });
看起来好像更简洁,但会出大问题——log() 函数的参数 message 是一个局部变量,log() 一旦返回,这个局部变量就被销毁了。而 strand 队列里的任务是异步执行的,很可能等真正执行到 do_log(message) 的时候,log() 早就返回、message 对应的内存早就被回收或者被别的数据覆盖了。结果就是:写到日志文件里的内容变得残缺不全,甚至出现乱码字符。
这提醒我们一个很重要的原则:在异步编程里,永远要假设"我以为立刻发生的事情,实际上会被延后执行"。操作系统、线程调度器什么时候真正跑到你的回调代码,不是你能控制的,你唯一能控制的是"传给回调的数据,在回调真正执行的那一刻是不是还活着"。所以涉及到异步回调时,尽量按值捕获(拷贝一份),或者用 shared_ptr 之类的手段显式延长对象生命周期。
6. 另一种写法:用 std::bind 代替 lambda
除了 lambda,也可以用 std::bind 达到同样的效果:
void log(const std::string message) {
boost::asio::post(strand_, std::bind(&Logger::do_log, this, message));
}
std::bind(&Logger::do_log, this, message) 的意思是:把成员函数 do_log、调用它所需要的对象指针 this、以及参数 message 打包成一个可调用对象。std::bind 内部同样会拷贝一份 message,所以生命周期问题和 lambda 按值捕获是一样安全的。现代 C++ 里 lambda 表达式通常更直观、更容易读,所以更推荐用 lambda,但了解 std::bind 这种写法在维护老代码时会用得上。
7. strand 到底在干什么:图解
7.1 用文字图理解 strand 的"排队漏斗"效果
多个线程并发调用 log() strand内部队列(先进先出) 唯一的 io_context 工作线程
------------------------ ------------------------ --------------------------
线程0: log("T0 msg0") --post--> [T0msg0][T1msg0][T2msg0][T3msg0]... --依次取出--> do_log() 一次只执行一个任务
线程1: log("T1 msg0") --post-->
线程2: log("T2 msg0") --post-->
线程3: log("T3 msg0") --post-->
任务只能一个接一个被取走执行,
绝不会有两个 do_log() 同时在跑
可以把 strand 想象成一个漏斗:不管上面有多少个线程同时往里面倒任务,漏斗口每次只能出来一个,任务在里面自动排好队,下面接水的人(负责真正执行任务的线程)永远只用应付"一次一个"的情况,完全不用担心冲突。
7.2 用时序图理解一次写日志的完整过程
8. 小结
- strand 是 Boost.Asio 提供的一种"不用锁也能保证串行执行"的机制,本质是给异步任务排一条队,保证同一时刻只有一个任务在跑。
- 隐式 strand 出现在两种情况:只用一个线程跑
io_context::run();或者异步操作在回调里链式发起下一个异步操作。 - 显式 strand(
boost::asio::strand)用于多线程场景:把需要串行化的操作(比如写同一个文件)统一通过post(strand, 任务)提交,strand 会自动帮你排好顺序。 - 用
work_guard防止io_context在还有后续任务要来的时候提前退出。 - lambda 捕获变量时,异步场景下优先按值捕获,避免引用悬空导致数据损坏;
std::bind效果类似,也是拷贝参数。 - strand 保证的是"不交叉、按post顺序执行",不代表各个来源的任务会被"公平轮流"处理,具体顺序取决于任务实际被提交的时间点。
本章总结:Boost.Asio 异步编程全景梳理
一、这一章到底在讲什么
这一章整体是在讲 Boost.Asio 这个库,核心目标是搞清楚:当程序需要和操作系统管理的外部资源(比如网络连接、定时器、文件描述符)打交道时,怎么用 Boost.Asio 去组织和管理这些异步任务。
可以把整章内容理解成一条主线:从"最基础的两个核心对象是什么" 讲到 “怎么把多个异步任务安全地组织起来”,再到"怎么控制这些任务的生命周期和执行方式",最后延伸到"和协程、网络编程结合起来怎么用"。下面按照这条主线把各个知识点串一遍。
二、核心基础:I/O 对象 与 I/O 执行上下文对象
这一章最先建立的两个基本概念:
- I/O 对象(I/O object):具体做事情的那个东西,比如一个定时器(
steady_timer)、一个 socket。它代表"某一种具体的异步资源"。 - I/O 执行上下文对象(I/O execution context,也就是
io_context):负责统一调度、驱动这些 I/O 对象去执行异步操作的"发动机"。所有的异步任务最终都要靠它的事件循环去真正跑起来。
这一章花了不少篇幅讲清楚这两者是怎么配合工作的: - I/O 对象怎么和操作系统的底层服务(比如系统调用、内核事件通知机制)打交道、互相通信;
- 这套设计背后遵循的设计原则是什么(比如"发起操作时立刻返回,不阻塞调用者,真正完成时再回调"这种异步模型);
- 在单线程程序和多线程程序里,应该怎么正确地使用它们,避免出现数据竞争或者对象生命周期上的问题。
三、任务的组织与控制:串行化、生命周期、启动与取消
在打好基础之后,这一章介绍了几种实际开发中非常关键的技巧:
- 用 strand 做工作的串行化:当多个异步操作可能在不同线程里并发执行、但又需要避免相互竞争同一份数据时,可以借助 strand 把这些操作按顺序排队执行,而不需要额外加锁。
- 管理异步操作中用到的对象的生命周期:异步操作是"发起后不会马上执行完"的,中间这段等待期间,参与操作的对象必须保证还"活着",这一章讲了怎么正确地管理这些对象,避免出现对象提前销毁、回调访问野指针这类问题。
- 任务的启动、中断与取消:包括怎么发起一个异步任务、怎么在任务执行到一半时打断它、以及怎么彻底取消一个还没完成的任务(比如前面提到的对象级取消和更精细的按操作取消)。
- 事件处理循环的管理:也就是
io_context内部那套不断轮询、分发已完成事件的循环机制,这一章讲了怎么控制这个循环的运行方式(比如什么时候启动、什么时候停止、要不要多线程一起跑这个循环)。 - 处理操作系统发来的信号:程序运行时操作系统可能会发送各种信号(比如用户按下
Ctrl+C),这一章也讲了 Boost.Asio 提供的信号处理方式,让程序能优雅地响应这些信号,而不是被动挨打。
四、延伸内容:网络编程与协程
除了上面这些核心机制,这一章还顺带引入了两个和异步编程紧密相关的话题:
- 网络编程相关的概念:把前面讲的这套异步框架,落地到实际的网络通信场景里去理解和使用。
- 协程(coroutine)相关的概念:协程是另一种组织异步代码的方式,能让异步逻辑写起来更接近"顺序执行"的直觉,这一章对这块内容做了初步的铺垫介绍。
并且,这一章通过若干实际例子,把上面讲到的这些概念都落实成了可以运行的代码,帮助从"知道原理"过渡到"会真正写出来"。
五、这一章带来的收获
学完这一章,主要收获可以归纳成两方面:
- 对"如何在 C++ 里管理异步任务"这件事,有了更系统、更深入的理解,不再是零散的知识点;
- 对 Boost.Asio 这个被广泛使用的库,理解了它底层到底是怎么实现异步调度这套机制的,而不只是停留在"会调用几个 API"的层面。
六、承上启下:下一章要讲什么
下一章会转向另一个 Boost 库——Boost.Cobalt。这个库提供了一套更高层、更丰富的接口,专门用来基于协程来开发异步程序。可以理解为:这一章讲的是异步编程"偏底层、偏原理"的那一套机制,而下一章要讲的是"建立在协程之上、用起来更顺手"的高层封装,两者是递进关系。
七、本章知识结构一览
八、章节内容结构(文本形式)
本章:Boost.Asio 异步编程
│
├── 核心基础
│ ├── I/O 对象(做具体事情的资源,如定时器、socket)
│ └── io_context(驱动异步任务执行的调度中心)
│
├── 任务组织与控制
│ ├── strand:串行化并发工作
│ ├── 生命周期管理:保证异步操作用到的对象不会提前失效
│ ├── 启动 / 中断 / 取消任务
│ ├── 事件处理循环的管理方式
│ └── 操作系统信号的处理
│
├── 延伸话题
│ ├── 网络编程相关概念
│ └── 协程相关概念
│
└── 实践
└── 若干可运行示例,把前面概念落地
openEuler 是由开放原子开源基金会孵化的全场景开源操作系统项目,面向数字基础设施四大核心场景(服务器、云计算、边缘计算、嵌入式),全面支持 ARM、x86、RISC-V、loongArch、PowerPC、SW-64 等多样性计算架构
更多推荐
所有评论(0)