1. 项目概述:为什么需要发送队列?
在之前的几篇关于asio网络编程的分享里,我们搭建了基础的客户端和服务器,实现了异步的读写操作。很多朋友跟着做下来,可能会发现一个“看似能用,实则暗藏隐患”的问题:当我们尝试在短时间内连续调用async_write发送多条数据时,程序的行为会变得不可预测,甚至直接崩溃。这背后的核心原因,是asio的异步写操作并非线程安全的,并且对并发调用有严格的限制。
简单来说,asio::async_write是一个“启动并遗忘”的操作。你调用它,它告诉操作系统“请把这些数据发出去”,然后就立刻返回了,不会阻塞你的线程。但是,如果你在前一个async_write操作还没完成(即操作系统内核还在处理发送缓冲区,或者网络拥塞导致数据没发完)的时候,立刻又启动一个新的async_write去发送另一段数据,那么这两段数据在底层套接字的发送缓冲区里就会“打架”,导致数据错乱、覆盖,最终引发程序崩溃。这就好比你在厨房用一个锅炒菜,菜还没盛出来,你又急着往同一个锅里倒进新的食材,结果可想而知。
所以,“发送队列”就是为了解决这个问题而生的。它的核心思想是:将“数据生产”和“网络发送”这两个动作解耦。所有需要发送的数据,都先放入一个队列(一个先进先出的容器)里。我们保证,在任何时刻,最多只有一个async_write操作在进行。只有当这个操作完成后,我们才从队列里取出下一个数据包,启动下一次发送。这样,无论上层业务逻辑以多快的频率产生数据,底层网络发送都能有条不紊、一个接一个地进行,从而真正实现稳定、可靠的全双工通信(即客户端和服务器可以同时、独立地收发数据)。
2. 核心设计思路与架构拆解
2.1 全双工通信的挑战与队列的角色
全双工通信意味着连接的两端都可以同时进行读和写。读操作(async_read)相对简单,因为通常我们用一个固定的缓冲区等待数据到来即可。但写操作是主动的、由我们触发的,其频率和时机不可控。
在没有队列的情况下,一个典型的错误模式是这样的:
- 用户点击按钮,触发发送消息A。
- 程序调用
async_write(A)。 - 在消息A还在发送的过程中,用户又快速点击了按钮,触发发送消息B。
- 程序在另一个线程(或主线程的事件循环中)直接调用
async_write(B)。 - BOOM!两个异步写操作同时操作底层套接字,导致未定义行为。
发送队列引入了一个“缓冲区”和“调度器”的角色:
- 缓冲区(队列本身):临时存储所有待发送的数据包。它可以是
std::deque、std::list或std::queue,包裹着我们要发送的数据(如std::string或std::vector<char>)。 - 调度器(发送逻辑):一个状态机,它只做两件事:
- 检查状态:当前是否有写操作正在进行?如果没有,且队列不为空,则进行步骤2。
- 执行发送:从队列头部取出一个数据包,启动一个
async_write操作。并为这个操作设置一个完成回调函数(Completion Handler)。
这个设计的关键在于,“启动发送”这个动作,永远只发生在两个时机:1) 队列从空变为非空时;2) 上一个写操作完成时。这就保证了串行化。
2.2 工具选型:为什么用std::deque而不用std::queue?
在C++标准库中,std::queue通常作为容器适配器,默认底层使用std::deque。两者都能满足先进先出的需求。但我个人更倾向于直接使用std::deque,原因有两点:
- 调试便利性:
std::deque支持迭代器,在调试时你可以直观地看到队列里所有排队的数据包内容,方便排查问题。std::queue的接口更为封闭。 - 内存分配考量:
std::deque通常由一系列固定大小的块(chunks)组成,在两端添加/删除元素效率都很高,且不会导致所有元素的大规模内存搬移。这对于一个可能频繁入队和出队的发送队列来说是很合适的特性。
当然,std::list也是一个选项,它的元素插入删除是常数时间,且指针稳定性最好。但std::list的内存开销(每个元素都需要额外的前后指针)和内存碎片化可能更严重。对于网络数据包这种“小对象但数量可能多”的场景,std::deque在内存局部性和综合性能上往往表现更好。
注意:这里的选择没有绝对的对错,取决于你的具体场景。如果数据包非常大(比如几MB),
std::list的指针稳定性优势会更明显。但对于常规的聊天消息、游戏指令等(几KB到几十KB),std::deque是更常见的选择。
2.3 线程安全与锁的选择
我们的网络IO操作(async_read,async_write)都是在asio的io_context事件循环所在的线程(通常称为IO线程)中发起和完成回调的。但是,将数据放入发送队列(post或push)这个动作,可能发生在任何线程。比如,你的UI线程收到用户输入,或者一个业务逻辑线程处理完数据后需要发送。
因此,对发送队列(std::deque)的访问(push_back和pop_front)必须是线程安全的。我们需要一把锁。
std::mutex:这是最直接的选择。在入队和检查队列状态时加锁。asio::strand:这是asio提供的一个更高级的抽象。strand可以确保所有通过它post或dispatch的函数对象(handler)都被序列化执行,即使它们来自不同的线程。你可以把strand理解为一个特殊的“序列化执行器”。
如何选择?
- 如果你的程序逻辑简单,所有可能操作队列的地方,你都方便拿到同一个
std::mutex,那么用mutex没问题。 - 如果你希望更紧密地与asio集成,并且你的异步操作链比较复杂,使用
asio::strand是更“asio风格”的做法。它可以保证所有相关的回调都在同一个逻辑线程上执行,无需显式加锁,避免了死锁风险。
在本篇的实现中,为了概念清晰,我们先使用std::mutex。但在一个更复杂的生产环境中,我会强烈建议使用asio::strand来管理所有与某个连接相关的异步操作(包括读回调、写回调、队列操作)。
3. 发送队列的详细实现步骤
3.1 定义连接类与数据结构
首先,我们定义一个TcpConnection类,它代表一个TCP连接,并内置发送队列功能。这里使用std::shared_ptr来管理连接的生命周期,这是asio网络编程中的常见模式。
// tcp_connection.hpp #ifndef TCP_CONNECTION_HPP #define TCP_CONNECTION_HPP #include <asio.hpp> #include <deque> #include <memory> #include <mutex> #include <string> using asio::ip::tcp; class TcpConnection : public std::enable_shared_from_this<TcpConnection> { public: using Pointer = std::shared_ptr<TcpConnection>; static Pointer Create(asio::io_context& io_context) { return Pointer(new TcpConnection(io_context)); } tcp::socket& Socket() { return socket_; } void Start(); // 开始读写 void Send(const std::string& message); // 供外部调用的发送接口 private: TcpConnection(asio::io_context& io_context); void DoRead(); // 执行异步读 void DoWrite(); // 执行异步写(从队列取数据) void OnWriteComplete(const asio::error_code& error, std::size_t bytes_transferred); // 写完成回调 tcp::socket socket_; std::array<char, 8192> read_buffer_; // 读缓冲区 // 发送队列相关成员 std::deque<std::string> write_queue_; // 发送队列 std::mutex queue_mutex_; // 保护队列的互斥锁 bool is_writing_; // 标志位:是否正在写入 }; #endif // TCP_CONNECTION_HPP关键成员解析:
write_queue_: 这就是我们的发送队列,存储待发送的字符串。queue_mutex_: 保护write_queue_和is_writing_标志的互斥锁。is_writing_: 一个非常重要的布尔标志。它表示当前是否有一个async_write操作正在进行中。绝对不要依赖write_queue_.empty()来判断是否正在写,因为异步操作是并发的。这个标志是保证串行化的核心。
3.2 实现核心的发送逻辑
让我们看看Send方法和DoWrite、OnWriteComplete是如何协作的。
// tcp_connection.cpp (部分) void TcpConnection::Send(const std::string& message) { // 1. 将数据包放入队列(需要加锁) { std::lock_guard<std::mutex> lock(queue_mutex_); write_queue_.push_back(message); } // 2. 尝试启动写操作 // 注意:这里不能直接调用DoWrite,因为要判断 is_writing_ 标志。 // 我们使用asio::post确保判断和启动写操作在同一个线程(IO线程)中执行,避免竞态条件。 asio::post(socket_.get_executor(), [self = shared_from_this()]() { // 捕获shared_ptr以延长连接生命周期 std::lock_guard<std::mutex> lock(self->queue_mutex_); // 如果当前没有正在进行的写操作,且队列里有数据,则启动写 if (!self->is_writing_ && !self->write_queue_.empty()) { self->is_writing_ = true; self->DoWrite(); // 启动实际的异步写 } // 否则(正在写),数据已经入队,等待当前写操作完成后的回调来处理下一个 }); }Send函数做了两件事:1) 安全地将数据入队;2) 通过asio::post将一个任务投递到IO线程,这个任务会检查is_writing_标志,如果空闲则启动写操作。
重要心得:为什么要在
asio::post的回调里加锁判断,而不是在Send函数里判断?因为is_writing_标志可能在OnWriteComplete回调中被修改,而这个回调也运行在IO线程。通过asio::post,我们确保了“检查标志”和“启动写”这两个动作与“完成回调修改标志”在同一个线程序列中执行,避免了复杂的跨线程同步问题。这是asio编程中保证线程安全的常用模式。
void TcpConnection::DoWrite() { // 这个函数总是在持有 queue_mutex_ 锁且 is_writing_ == true 的情况下被调用 if (write_queue_.empty()) { // 防御性编程:理论上不会进入这里,但如果发生,需要重置状态 std::lock_guard<std::mutex> lock(queue_mutex_); is_writing_ = false; return; } // 取出队列头部的数据包 const std::string& packet_to_send = write_queue_.front(); // 发起异步写操作 asio::async_write(socket_, asio::buffer(packet_to_send.data(), packet_to_send.size()), [self = shared_from_this()](const asio::error_code& ec, std::size_t bytes_transferred) { // 写操作完成,回调到OnWriteComplete self->OnWriteComplete(ec, bytes_transferred); }); }DoWrite函数假设它被调用时,队列非空且is_writing_为真。它取出队首数据,启动异步写。
void TcpConnection::OnWriteComplete(const asio::error_code& error, std::size_t bytes_transferred) { if (error) { // 发生错误:连接可能已断开 std::cerr << "Write failed: " << error.message() << std::endl; // 处理错误,例如关闭socket return; } // 写成功,移除已发送的数据包 { std::lock_guard<std::mutex> lock(queue_mutex_); write_queue_.pop_front(); // 移除已发送的包 // 检查队列是否还有数据 if (write_queue_.empty()) { // 队列已空,停止写循环 is_writing_ = false; } else { // 队列还有数据,继续发送下一个包 // is_writing_ 保持为 true DoWrite(); // 递归调用,发送下一个 } } }OnWriteComplete是逻辑的核心:
- 处理错误。
- 成功则移除已发送的包。
- 如果队列变空,则重置
is_writing_标志,写循环停止。 - 如果队列还有数据,则递归调用
DoWrite(),发送下一个包。这里形成了一个“链式调用”,一个接一个地发送,直到队列清空。
3.3 启动连接与读操作
读操作相对独立,与发送队列无关,但为了完整性,这里给出Start和DoRead的实现。
void TcpConnection::Start() { DoRead(); // 开始读循环 // 注意:这里不自动启动写。写操作由外部调用Send触发。 } void TcpConnection::DoRead() { auto self(shared_from_this()); socket_.async_read_some(asio::buffer(read_buffer_), [this, self](const asio::error_code& ec, std::size_t length) { if (!ec) { // 处理读到的数据,例如打印或转发 std::string received_data(read_buffer_.data(), length); std::cout << "Received: " << received_data << std::endl; // 可以在这里触发业务逻辑... // 继续读 DoRead(); } else { // 读错误,连接关闭或出错 std::cerr << "Read error: " << ec.message() << std::endl; // 清理资源... } }); }4. 服务器与客户端的集成示例
4.1 服务器端实现
服务器端使用我们刚实现的TcpConnection类。
// tcp_server.cpp #include "tcp_connection.hpp" #include <asio.hpp> #include <iostream> #include <set> class TcpServer { public: TcpServer(asio::io_context& io_context, short port) : acceptor_(io_context, tcp::endpoint(tcp::v4(), port)) { DoAccept(); } private: void DoAccept() { acceptor_.async_accept( [this](const asio::error_code& ec, tcp::socket socket) { if (!ec) { auto conn = TcpConnection::Create(socket.get_executor().context()); conn->Socket() = std::move(socket); connections_.insert(conn); conn->Start(); std::cout << "New connection accepted. Total: " << connections_.size() << std::endl; // 示例:向新连接发送欢迎消息 conn->Send("Welcome to the server!\n"); } else { std::cerr << "Accept error: " << ec.message() << std::endl; } // 继续接受新连接 DoAccept(); }); } tcp::acceptor acceptor_; std::set<std::shared_ptr<TcpConnection>> connections_; // 管理所有活跃连接 }; int main() { try { asio::io_context io_context; TcpServer server(io_context, 12345); std::cout << "Server started on port 12345" << std::endl; io_context.run(); // 启动事件循环 } catch (std::exception& e) { std::cerr << "Exception: " << e.what() << std::endl; } return 0; }4.2 客户端实现
客户端同样使用TcpConnection,并模拟快速连续发送。
// tcp_client.cpp #include "tcp_connection.hpp" #include <asio.hpp> #include <iostream> #include <thread> #include <chrono> int main() { try { asio::io_context io_context; // 解析服务器地址 tcp::resolver resolver(io_context); auto endpoints = resolver.resolve("127.0.0.1", "12345"); // 创建连接 auto connection = TcpConnection::Create(io_context); // 异步连接 asio::async_connect(connection->Socket(), endpoints, [connection](const asio::error_code& ec, const tcp::endpoint&) { if (!ec) { std::cout << "Connected to server!" << std::endl; connection->Start(); // **模拟快速连续发送,测试队列** std::cout << "Sending 10 messages rapidly..." << std::endl; for (int i = 0; i < 10; ++i) { connection->Send("Message " + std::to_string(i) + "\n"); // 不加延时,瞬间发送 } // 再发送一个稍大的消息 connection->Send(std::string(1000, 'X')); // 1000个'X' } else { std::cerr << "Connect failed: " << ec.message() << std::endl; } }); // 在另一个线程中运行io_context,以便主线程可以做其他事(例如接收用户输入) std::thread io_thread([&io_context]() { io_context.run(); }); // 主线程:模拟用户输入发送 std::string user_input; while (std::getline(std::cin, user_input)) { if (user_input == "quit") break; // 这里需要注意:connection是在io_context的线程中使用的。 // 我们需要通过post将Send操作投递到io_context的线程中执行。 asio::post(io_context, [connection, user_input]() { connection->Send("[Client says]: " + user_input + "\n"); }); } io_context.stop(); io_thread.join(); } catch (std::exception& e) { std::cerr << "Exception: " << e.what() << std::endl; } return 0; }5. 常见问题、性能考量与进阶优化
5.1 为什么我的程序在发送大量数据时内存暴涨?
这是实现发送队列时最容易掉进去的坑。看我们之前的实现,队列里存储的是std::string,也就是数据的副本。如果你要发送一个1MB的字符串,队列里就存着一个1MB的std::string。如果网络很慢(比如发送速度是100KB/s),而你的生产速度很快(比如每秒产生10个1MB的数据包),那么队列会迅速堆积,导致内存耗尽。
解决方案:使用std::shared_ptr<const std::string>或者自定义的缓冲区对象。
// 在连接类中 std::deque<std::shared_ptr<const std::string>> write_queue_; void TcpConnection::Send(const std::string& message) { auto packet = std::make_shared<std::string>(message); // 在堆上分配,引用计数 { std::lock_guard<std::mutex> lock(queue_mutex_); write_queue_.push_back(packet); } // ... 后续投递逻辑不变 } void TcpConnection::DoWrite() { auto packet_to_send = write_queue_.front(); // 取出的是 shared_ptr asio::async_write(socket_, asio::buffer(packet_to_send->data(), packet_to_send->size()), [self = shared_from_this(), packet_to_send](const asio::error_code& ec, std::size_t bytes_transferred) { // 注意:回调里捕获了 packet_to_send,这意味着在异步操作完成前, // 这个 shared_ptr 会一直保持数据存活,防止数据被提前销毁。 // 操作完成后,packet_to_send 离开作用域,引用计数减一。 self->OnWriteComplete(ec, bytes_transferred); }); }这样做的好处是,数据在堆上只有一份,队列里存的只是指向它的智能指针,拷贝成本很低。更重要的是,在异步写操作进行时,回调函数通过捕获shared_ptr保持了数据的生命,发送完成后自动释放。队列里即使有多个指向同一份大数据的指针,内存占用也只是一份。
5.2 如何实现发送流量控制?
发送队列解决了并发调用的问题,但如果接收方处理速度远慢于发送方,队列还是会无限增长。这就需要更高级的流量控制(背压,Backpressure)。
一个简单的方法是在连接对象中增加一个“高水位线”(High Water Mark)。
class TcpConnection { // ... static const size_t MAX_QUEUE_SIZE = 100; // 最大队列长度 std::atomic<size_t> pending_write_bytes_{0}; // 待发送总字节数(近似值) // ... void Send(const std::string& message) { size_t packet_size = message.size(); // 检查是否超过高水位线 if (pending_write_bytes_.load() + packet_size > MAX_QUEUE_SIZE * 1024) { // 假设以KB为单位 std::cerr << "Send queue full, dropping packet." << std::endl; // 可以选择丢弃、阻塞或返回错误给上层 return; } pending_write_bytes_ += packet_size; // ... 入队逻辑 } void OnWriteComplete(...) { // ... { std::lock_guard<std::mutex> lock(queue_mutex_); auto sent_packet = write_queue_.front(); pending_write_bytes_ -= sent_packet->size(); // 发送完成,减去字节数 write_queue_.pop_front(); // ... } // ... } };更复杂的流量控制需要应用层协议支持,例如TCP本身的滑动窗口是传输层的流量控制,而应用层可以定义类似“ACK”或“READY”的信号,让接收方告诉发送方“我准备好了,你可以再发N个数据包”。
5.3 使用asio::strand替代std::mutex
如前所述,使用strand是更优雅的方式。修改如下:
class TcpConnection { // ... private: asio::strand<asio::io_context::executor_type> strand_; // 增加strand std::deque<...> write_queue_; // 移除 queue_mutex_ 和 is_writing_ bool is_writing_; // 这个标志现在不需要了,因为strand保证了串行访问 }; TcpConnection::TcpConnection(asio::io_context& io_context) : socket_(io_context), strand_(asio::make_strand(io_context)) { // 初始化strand } void TcpConnection::Send(const std::string& message) { // 通过strand.post,确保入队和后续检查在同一个逻辑线程序列中执行 asio::post(strand_, [self = shared_from_this(), message]() { // 这里直接捕获message副本 self->write_queue_.push_back(message); // 无需加锁 // 如果队列大小为1(即之前是空的),说明需要启动写操作 if (self->write_queue_.size() == 1) { self->DoWrite(); } // 如果队列大小>1,说明已经有写操作在进行(DoWrite链正在运行), // 新数据入队后会自动被后续的DoWrite处理。 }); } void TcpConnection::DoWrite() { // 这个函数总是在strand中执行,所以访问write_queue_是安全的 if (write_queue_.empty()) { return; // 队列已空,停止链 } const std::string& packet = write_queue_.front(); asio::async_write(socket_, asio::buffer(packet), asio::bind_executor(strand_, // **关键**:确保完成回调也在同一个strand中执行 [self = shared_from_this()](const asio::error_code& ec, std::size_t bytes) { // 这个回调在strand中,所以访问write_queue_安全 if (!ec) { self->write_queue_.pop_front(); // 递归调用,发送下一个 self->DoWrite(); } else { // 错误处理 std::cerr << "Write error in strand: " << ec.message() << std::endl; } })); }使用strand后,所有对write_queue_的访问(push_back,pop_front,front,empty)都因为被post或bind_executor到同一个strand而自动序列化,完全消除了显式锁的需要,代码更简洁,更不容易死锁。
5.4 性能瓶颈与优化方向
- 队列数据结构:对于超高性能场景,
std::deque可能因为内存块分配和缓存不友好成为瓶颈。可以考虑使用无锁队列(如moodycamel::ConcurrentQueue),但复杂度会急剧上升。绝大多数情况下,std::deque或std::list配合strand足够了。 - 内存分配:频繁构造/析构
std::string或shared_ptr会产生内存分配开销。可以使用内存池或对象池来复用缓冲区对象。 - 小包合并(Nagle算法与TCP_CORK):TCP有Nagle算法来合并小包,但有时为了低延迟需要禁用它(
TCP_NODELAY)。另一方面,如果你发送的是大量小消息,可以在应用层做一个简单的合并:当数据包很小且队列中已有数据时,不立即启动异步写,而是等待一个极短的时间(如1ms)或积累到一定大小(如4KB)后再一次性发送。这需要更精细的定时器控制。 - 异步写回调的代价:每个数据包都对应一个异步操作和一次回调。如果每秒要发送数万个极小的数据包,回调开销可能显著。这时可以考虑在
DoWrite中,一次性将队列中的多个连续小包通过asio::write(同步)或构造一个asio::const_buffer序列传给async_write,减少回调次数。但这会提高实现的复杂性。
实现一个健壮的发送队列,是构建高性能、高可靠性网络服务的基石。它看似只是加了个容器,实则涉及线程安全、生命周期管理、流量控制、性能优化等多个核心知识点。从最简单的mutex+deque开始,理解其工作原理,再逐步演进到使用strand、智能指针管理缓冲区、乃至实现背压机制,这个过程本身就是对异步编程和网络编程思想的深度锤炼。