ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Boost.Asio实战:异步网络编程、Strand并发控制与批量写库

Boost.Asio实战:异步网络编程、Strand并发控制与批量写库 去年接手一个物联网网关项目要把上千台终端设备的实时数据接进来再写进时序数据库。网络层选型时我几乎没有犹豫就定了Boost.Asio。不是说裸socket和epoll不行而是连接数量上去之后手动维护状态机、处理半包、跨平台迁移每一步都在消耗精力。Asio把这些琐碎但关键的逻辑统一抽象成异步操作回调让业务逻辑可以集中在一个on_read/on_write里。这篇文章不打算复述官方文档只讲我自己实际跑过、坑过、优化过的经验从同步server起步讲到异步模型、回调生命周期和strand再讲如何把收到的数据安全地批量写入数据库。不管你是刚接触C网络编程还是已经写过一阵子socket但没上过Asio应该都能找到有用的东西。1. 为什么最后选了Boost.Asio手写Socket的痛Asio刚好都能治1.1 每连接一个线程为什么扛不住早期做局域网小工具用socket加线程一个连接一个std::thread连接少的时候挺好。到了网关这种要支撑几千路长连接的项目这个模型直接暴雷。首先是线程切换成本。两千个线程每个栈默认8MB虚拟内存光栈空间就吃掉不少再加上上下文切换CPU时间大片浪费在线程调度上。其次是同步阻塞带来的连锁反应某个连接上的业务如果慢了对应线程被卡住Accept侧还得继续拉线程数量一路暴涨。后来转epoll属于另一种痛苦。边沿触发还是水平触发、事件表怎么维护、连接断开时状态怎么清理、业务层协议怎么从字节流里切包全部自己造。写出来能跑但代码量很大而且换到Windows又要重来。我一度转向libevent可C接口的复杂度不低回调里到处是void*上下文C这边维护起来还是别扭。Boost.Asio的抽象从根本上解决了这些问题。1.2 Asio的核心抽象IO对象加Proactor模式很多人以为Asio只是把epoll包了一层其实它借鉴的是Proactor模式你发起一个async_read给它一个完成回调底层等待操作完成然后把回调投递到io_context上执行。用户看到的不是这个socket可读了而是读完了数据在这里。这个差异非常关键。epoll给你的是事件到达通知你还要自己去判断读多少、是否读完Asio直接把操作完成这个结果交给你的业务代码。Linux上它用epoll实现Windows上落到IOCP接口完全一致。跨平台这点对需要适配多种部署环境的服务来说是实打实的省力。io_context是整个库的核心调度器它自己不会主动干活必须调用run()才会投产。可以把run()理解为进入事件循环不断取出已完成操作的handler并执行。这个设计让线程模型非常灵活你可以单线程跑也可以多个线程同时调run()让handler并行执行。1.3 环境搭建头文件、链接库、第一个编译问题起步阶段卡人的往往不是概念而是编译链接。Ubuntu/Debian下sudo apt install libboost-all-dev g -stdc17 -O2 -pthread main.cpp -lboost_system -o serverBoost 1.66之后Asio主体基本是header-only不需要单独链接libboost_system但如果你用到boost::system::error_code或者使用的是较老版本Boost依然要-lboost_system。与其省这个链接我建议一开始就加上避免老项目升级时踩undefined reference的坑。Windows上我一般用vcpkg安装vcpkg install boost-asio boost-system然后集成到CMake注意只保留一份Boost别把系统里的和vcpkg里的混用符号版本冲突排查起来很麻烦。这里提一个常见报错编译通过但链接时报undefined reference to boost::system::detail::...多半就是缺了boost_system或者Boost头文件与库版本不一致。2. 先把同步模型跑透Echo Server的骨架和buffer的真实语义2.1 十几行代码搭一个同步TCP服务端网上有很多博客上来直接甩异步代码读者看完云里雾里。我建议先跑通同步模型把acceptor、socket、buffer这几个基础概念落到实处。下面是一个最简同步echo server#include boost/asio.hpp #include iostream using boost::asio::ip::tcp; int main() { try { boost::asio::io_context io; tcp::acceptor acceptor(io, tcp::endpoint(tcp::v4(), 8080)); std::cout listen on 8080 std::endl; while (true) { tcp::socket socket(io); acceptor.accept(socket); std::arraychar, 128 buf; boost::system::error_code ec; size_t len socket.read_some(boost::asio::buffer(buf), ec); if (ec) { std::cerr read failed: ec.message() std::endl; continue; } boost::asio::write(socket, boost::asio::buffer(buf, len), ec); if (ec) { std::cerr write failed: ec.message() std::endl; } } } catch (std::exception e) { std::cerr exception: e.what() std::endl; } return 0; }这段代码的关键在于acceptor负责监听accept阻塞到有连接进来read_some读一段数据write原样回写。异常方面boost::asio多数的函数都有抛异常版本和返回error_code版本服务端代码我更习惯用error_code版本把业务逻辑里面的正常关闭、异常断开分开处理。2.2 buffer是视图不是容器它不拷贝只是描述内存区间boost::asio::buffer是新手容易理解错的地方。它不是容器不拥有内存只是包装了起始地址长度的一个视图。std::string s hello; boost::asio::write(sock, boost::asio::buffer(s)); // 同步写没问题同步操作里write返回时数据已经写入OS发送缓冲区内存安全性比较简单。但如果是异步操作这里就暗藏危机void bad_write(tcp::socket sock, boost::asio::io_context io) { std::string s hello; boost::asio::async_write(sock, boost::asio::buffer(s), [](auto, size_t) {}); // 函数结束s 析构但异步写入可能还没触发回调 }异步回调触发时buffer指向的内存早就释放了轻则读到垃圾数据重则崩溃。在异步场景里所有传给async_*的buffer其底层内存必须活到回调执行完。2.3 同步模型的致命伤一个连接阻塞整个循环上面的server一次只能服务一个连接。只要某个客户端连上之后不发数据accept之后的read_some就卡在那里后面的连接全部排队。改成每个连接开一个线程while (true) { auto sock std::make_sharedtcp::socket(io); acceptor.accept(*sock); std::thread([sock]() { handle_client(*sock); }).detach(); }能并发但回到了1.1的问题。线程数随连接数上涨资源利用率低。所以同步模型适合写客户端、压测脚本、一次性工具真正的高并发服务还是要走异步。3. 异步模型进阶回调、生命周期和strand3.1 async_accept与enable_shared_from_this保命异步服务端经典写法是acceptor不断async_accept每次连接创建一个SessionSession用shared_ptr管理。class Session : public std::enable_shared_from_thisSession { public: explicit Session(tcp::socket socket) : socket_(std::move(socket)) {} void start() { do_read(); } private: void do_read() { auto self shared_from_this(); socket_.async_read_some(boost::asio::buffer(data_, max_length), [this, self](boost::system::error_code ec, std::size_t length) { if (!ec) { do_write(length); } }); } void do_write(std::size_t length) { auto self shared_from_this(); boost::asio::async_write(socket_, boost::asio::buffer(data_, length), [this, self](boost::system::error_code ec, std::size_t /*length*/) { if (!ec) { do_read(); } }); } tcp::socket socket_; enum { max_length 1024 }; char data_[max_length]; }; class Server { public: Server(boost::asio::io_context io, short port) : acceptor_(io, tcp::endpoint(tcp::v4(), port)) { do_accept(); } private: void do_accept() { acceptor_.async_accept( [this](boost::system::error_code ec, tcp::socket socket) { if (!ec) { std::make_sharedSession(std::move(socket))-start(); } do_accept(); }); } tcp::acceptor acceptor_; };为什么Session必须继承enable_shared_from_this因为异步链是一个环回调里捕获了Session对象回调执行后又注册下一个回调等于Session被回调链持有必须保证最后一个回调执行完之前对象不会析构。如果用裸指针或者栈对象回调触发时对象可能已经没了这是异步网络编程里最常见的内存问题。每次在do_read和do_write入口写auto self shared_from_this();是为了把引用计数先加一保证整个操作期间Session存活。然后在lambda捕获列表里带上self让生命周期顺着回调链一路延续。3.2 async_read_some与async_readTCP流没有消息边界async_read_some只保证读到一些字节不保证正好是一个业务消息。很可能一条业务消息拆成了两次到达也很可能两条消息粘在一次到达里。TCP是字节流消息边界必须由应用层自己定义。我习惯的做法是约定一个简单帧格式头部4字节表示长度后面跟着payload。然后在Session里维护一个接收缓冲std::vectorchar read_buf_; std::vectorchar frame_buf_;do_read不断把新数据追加到read_buf_然后尝试从read_buf_里切出完整帧void do_read() { auto self shared_from_this(); boost::asio::async_read(socket_, boost::asio::buffer(temp_buf_, temp_buf_.size()), [this, self](boost::system::error_code ec, std::size_t len) { if (ec) { handle_ec(ec); return; } read_buf_.insert(read_buf_.end(), temp_buf_.begin(), temp_buf_.begin() len); while (try_parse_one_frame()) {} do_read(); }); } bool try_parse_one_frame() { if (read_buf_.size() 4) return false; uint32_t body_len ntohl(*(uint32_t*)read_buf_.data()); if (read_buf_.size() 4 body_len) return false; // 取出完整包交给业务层 handle_frame(read_buf_.data() 4, body_len); read_buf_.erase(read_buf_.begin(), read_buf_.begin() 4 body_len); return true; }注意ntohl那行涉及字节序跨平台时要小心。更稳妥是用memcpy读长度避免未对齐访问。3.3 多线程run()之后strand是必需品io_context::run()可以同时在多个线程里调用这样多个handler会并发执行。但同一个socket上的读写如果并发就乱套了两个线程同时写数据交叉一个在读一个在关资源释放的时机也会出问题。strand就是解决这个问题的。它保证投递到同一个strand上的handler不会并发执行相当于给异步操作上了一把逻辑锁。boost::asio::strandboost::asio::io_context::executor_type strand_ boost::asio::make_strand(io); boost::asio::async_write(socket_, buffer, boost::asio::bind_executor(strand_, [this, self](boost::system::error_code ec, std::size_t len) { ... }));把Session里所有handler都通过bind_executor绑定到同一个strand这样即便io_context被多个线程跑同一Session内部仍然是串行的。新版Asio也可以直接用make_strand风格更简洁。很多人问能不能用一个全局mutex替代strand可以但代价是两个不同Session之间本来可以并行的读写也被串行化了吞吐直接下降。strand是细粒度的比全局锁合理得多。3.4 超时、心跳与deadline_timer搭配网络层最常见的故障就是半开连接对端已经消失本端还不知道。做法是用steady_timer配合读写操作做超时控制。boost::asio::steady_timer timer(io); timer.expires_after(std::chrono::seconds(30)); auto self shared_from_this(); timer.async_wait([this, self](const boost::system::error_code ec) { if (!ec) { // 超时主动断开 boost::system::error_code ignored; socket_.close(ignored); } });注意几点timer的回调也要绑定到同一个strand否则它可能和正在执行的读写回调并发另外每次读到数据后要把timer重置一下改成空闲超时而不是绝对超时。我一般把timer和读写操作都放在Session内部统一由strand串行化避免竞态。4. 打通业务系统把收到的数据批量写进数据库4.1 网络回调里写库是阻塞罪网关场景下服务端收到数据后通常要落库比如写进TDengine、MySQL、PostgreSQL。最容易犯的错误是在Asio的回调里直接调用数据库同步接口。io_context的线程就那么几个回调里一旦发生磁盘I/O和SQL执行这个线程就被占住了。执行时间一长后面所有连接的回调都排队网络延迟跟着飙升。数据库写慢100毫秒就可能拖累几十个连接。4.2 双缓冲队列加独立写线程把网络层与存储层解耦我采用的通用模式是网络回调只做解析和入队独立的工作线程批量取数写库。struct Record { uint64_t ts; double value; }; std::mutex mtx; std::dequeRecord queue; bool running true; void on_message(Record rec) { std::lock_guardstd::mutex lk(mtx); queue.push_back(std::move(rec)); } void db_worker() { while (running) { std::vectorRecord batch; { std::lock_guardstd::mutex lk(mtx); batch.assign(queue.begin(), queue.end()); queue.clear(); } if (batch.empty()) { std::this_thread::sleep_for(std::chrono::milliseconds(10)); continue; } batch_insert_to_db(batch); } }这个设计的好处是网络回调耗时可预测不会因为数据库抖动拖垮整个网络层同时写库是批量操作吞吐比单条插入高一个量级。4.3 接到TDengine场景taos_stmt_prepare批量绑定如果存储端是TDengine我给的方案是预编译语句批量绑定。TDengine的C接口提供了taos_stmt_prepare类似MySQL的prepared statement可以先把SQL模板准备好再用参数数组循环绑定。taos_stmt* stmt taos_stmt_init(taos_conn); const char* sql INSERT INTO ? USING metrics TAGS(?, ?) VALUES (?, ?); taos_stmt_prepare(stmt, sql, strlen(sql)); // 绑定表名和标签 taos_bind_t tb[1]; // ... 初始化每个字段的buffer、length、type taos_stmt_bind_param(stmt, tb, 1); // 追加多行值 taos_bind_t values[3]; // ... 初始化时间戳、设备ID、数值字段 for (size_t i 0; i batch.size(); i) { // 逐个绑定参数 values[0].u.var.i64 batch[i].ts; values[1].u.var.buflen ...; // 这里注意每个字段的buffer和len都要指向有效内存 taos_stmt_bind_param(stmt, values, 3); taos_stmt_add_batch(stmt); } taos_stmt_execute(stmt); taos_stmt_close(stmt);这段代码的重点是每次taos_stmt_bind_param传入的参数在调用期间必须有效add_batch会把当前参数复制进内部缓冲区所以可以复用values数组但字符串类型的字段要保证指针指向的内容不被提前释放。实际开发中不同TDengine版本的taos_bind_t字段名略有差异建议打开头文件确认一下或者用官方示例对照。还有一点批量大小不要贪多10到100条一批实测吞吐和内存占用都比较合适。这个思路同样适用于MySQL的prepare绑定接口原理是一样的。4.4 更现代的写法C20协程让异步代码变同步Boost.Asio从1.74开始积极拥抱C20协程用use_awaitable可以把回调嵌套改写得跟同步一样直观boost::asio::awaitablevoid handle_session(tcp::socket sock) { try { std::arraychar, 1024 buf; while (true) { size_t n co_await sock.async_read_some(boost::asio::buffer(buf), boost::asio::use_awaitable); co_await boost::asio::async_write(sock, boost::asio::buffer(buf, n), boost::asio::use_awaitable); } } catch (const boost::system::system_error e) { // 处理异常 } } boost::asio::awaitablevoid listener() { auto exec co_await boost::asio::this_coro::executor; tcp::acceptor acceptor(exec, {tcp::v4(), 8080}); while (true) { auto sock co_await acceptor.async_accept(boost::asio::use_awaitable); boost::asio::co_spawn(exec, handle_session(std::move(sock)), boost::asio::detached); } }协程底层仍然是handler好处是写复杂的收发顺序时不再一个套一个回调。但协程不是银弹协程栈上的对象生命周期要心里有数取消操作和超时也要显式处理。如果团队对C20协程不熟我建议先保持传统回调风格稳定再说。5. 这些年的错题本连接、内存、线程与性能5.1 对端关闭连接时EOF不是错误很多人在async_read_some回调里看到ec就当作错误处理打印日志、断开连接。实际上boost::asio::error::eof表示对端正常关闭了连接读完所有数据后收到FIN这不是异常是正常的业务结束。正确做法是把这个分支当作清理Session资源的信号而不是报警。真正需要注意的是connection_reset_by_peer这种异常断开比如对端进程崩溃、网络超时重扔RST。TCP这种场景很多不能一看到ec就panic式地把整个服务打挂要把可预期断开和异常断开分开处理。5.2 per-Session buffer固定长度可能浪费也容易踩踏我早期偷懒给每个Session固定一个4KB的数组当接收缓冲。后来发现有的消息只有几十字节有的消息有几十KB固定长度要么浪费内存要么一条大消息直接撑爆缓冲。后来改成三层结构小消息直接走栈上的临时缓冲区中等消息走Session内部vector大消息走单独的内存块。具体大小根据业务包特征调。其实核心原则很简单不要让所有连接都按最大包分配内存也不要让多个连接共享同一个缓冲区。共享缓冲区在多线程回调下就是定时炸弹。5.3 一个io_context配几个线程合适这是个老问题。我的经验值是纯转发场景CPU核数和io_context线程数差不多就行比如8核就开8个线程如果handler里有轻量计算但数据库操作放在单独的worker线程里那么网络线程开2~4个就够了开多了反而增加上下文切换如果handler里有耗时CPU计算比如解析超大JSON、协议编解码这部分最好挪到专用线程池或者增加网络线程数到CPU核数的1.5倍左右并加strand保护有一个容易忽略的点线程数超过CPU核数后靠并行提高吞吐的收益快速递减反而线程切换成了瓶颈。我见过有人给8核机器开了64个run()线程性能没涨多少CPU上下文切换却高了十倍。5.4 handler里无意拷贝导致内存疯涨Session写法规范后内存问题多半出在handler捕获里。比如这样std::string big_json; big_json ...; socket_.async_write(socket_, boost::asio::buffer(big_json), [this, self, big_json](auto...) {});这里big_json被按值捕获又复制了一份。如果消息量大等于每个发送队列里都有一份大字符串的副本内存很快飙上去。正确做法是让big_json的生命周期和异步操作绑定但不要无谓拷贝可以用std::shared_ptrstd::string把同一个对象传入buffer方法和lambda捕获。auto data std::make_sharedstd::string(big_json); boost::asio::async_write(socket_, boost::asio::buffer(*data), [this, self, data](auto...) {});这样底层buffer指向的还是data持有的那一段内存lambda再持有一份shared_ptr没有第二份数据副本。注意若数据本身不需要跨线程移动shared_ptr的引用计数更新虽有原子操作开销但通常远小于一次大内存拷贝。5.5 排错清单从现象到原因的快速对照现象可能原因处理方向编译链接报boost::systemundefined缺-lboost_system或Boost版本混乱检查路径只保留一份Boost服务启动报bind失败端口被占用或之前进程没退干净netstat -tunlp查端口等待TIME_WAIT释放或改用SO_REUSEADDR回调里访问Session崩溃没有用shared_from_this延长生命周期检查是否用裸指针捕获Session多条连接数据互相穿插共享了同一个buffer改成per-Session独立buffer高并发下CPU高但吞吐低线程数超过CPU核数频繁切换减少run()线程数检查是否有核间迁移内存随时间缓慢上涨handler按值拷贝大对象用shared_ptr持有数据避免二次拷贝我自己的排错习惯是先用strace看一下系统调用确认是网络层问题还是业务层问题再看io_context线程数是否合理最后才怀疑业务代码。网络库踩坑的规律往往是生命周期和资源所有权没想清楚而不是Asio本身不够好用。写到这里最想强调的还是那句老话异步网络编程也好Boost.Asio也好真正的难点并不在于某个API的用法而在于你愿不愿意把数据何时到达、对象何时销毁、回调何时并发这三件事想透。把这些基本功做扎实了不管前端的消费者是连接池、数据库还是协程风格的业务代码都能稳稳接住。
返回列表