Boost.Asio中async_write仅在服务器关闭后发送消息的问题求助
问题描述
我正在开发一款聊天服务器,能成功接收用户发送的消息,但调用async_write广播消息时,消息仅在服务器关闭(按下Ctrl+C)后才会发送至客户端。
示例:客户端发送“test”和“test2”,仅在关闭服务器后才收到拼接后的“testtest2”。
相关代码实现
消息发送代码
void Server::writeHandler(int id, boost::system::error_code error){ if (!error){ std::cout << "[DEBUG] message broadcasted "; } else { close_connection(id); } } void Server::broadcast(std::string msg, boost::system::error_code error){ for (auto& user : m_users){ // for every user in unordered map of users asio::async_write(user.second->socket, asio::buffer(msg, msg.size()), std::bind(&Server::writeHandler, this, user.first, std::placeholders::_1)); } }
onMessage中的broadcast调用
void Server::onMessage(int id, boost::system::error_code error){ if (!error){ broadcast(m_read_msg, error); // char m_read_msg[PACK_SIZE] // PACK_SIZE = 512 asio::async_read(m_users[id].get()->socket, asio::buffer(m_read_msg, PACK_SIZE), // PACK_SIZE = 512 std::bind(&Server::onMessage, this, id, std::placeholders::_1)); } else { close_connection(id); } }
服务器运行函数
void Server::run(int port){ asio::ip::tcp::endpoint endpoint(asio::ip::tcp::v4(), port); std::cout << "[DEBUG] binded on " << port << std::endl; // acceptor initialization m_acceptor = std::make_shared<asio::ip::tcp::acceptor>(m_io, endpoint); // start to listen for connections listen(); // run the io_service in thread io_thread = std::thread( [&]{ m_io.run(); } ); // m_io = io_serivce while (true){ } } void Server::listen(){ std::shared_ptr<User> pUser(new User(m_io)); m_acceptor->async_accept(pUser->socket, std::bind(&Server::onAccept, this, pUser)); }
解决方案
1. 修复异步写缓冲区的生命周期问题
asio::async_write要求缓冲区在整个异步操作期间保持有效。当前broadcast函数中的std::string msg是局部变量,函数返回后内存会被释放,导致异步写操作访问无效内存,引发未定义行为(比如消息被覆盖、延迟发送)。
修改方法:使用std::shared_ptr持有消息副本,确保在异步写完成前内存不会被销毁:
void Server::broadcast(std::string msg, boost::system::error_code error){ auto msg_ptr = std::make_shared<std::string>(std::move(msg)); for (auto& user : m_users){ asio::async_write(user.second->socket, asio::buffer(*msg_ptr), [this, user_id = user.first, msg_ptr](const boost::system::error_code& ec) { writeHandler(user_id, ec); }); } }
2. 避免共享读缓冲区的竞争
m_read_msg是全局成员变量,多个客户端的async_read会同时写入该缓冲区,导致消息内容被覆盖。应给每个客户端单独分配读缓冲区:
修改User类
class User { public: User(asio::io_service& io) : socket(io) {} asio::ip::tcp::socket socket; char read_buf[PACK_SIZE]; // 每个用户单独的读缓冲区,PACK_SIZE = 512 };
修改onMessage函数
void Server::onMessage(int id, boost::system::error_code error){ if (!error){ auto& user = m_users[id]; // 从用户专属缓冲区读取消息,去除末尾空字符 std::string msg(user->read_buf, PACK_SIZE); msg.erase(std::find(msg.begin(), msg.end(), '\0'), msg.end()); broadcast(msg, error); // 使用用户自己的缓冲区发起下一次异步读 asio::async_read(user->socket, asio::buffer(user->read_buf, PACK_SIZE), std::bind(&Server::onMessage, this, id, std::placeholders::_1)); } else { close_connection(id); } }
3. 优化主线程阻塞方式
原代码中的while(true)空循环会占用大量CPU资源,影响io线程的调度。可以使用休眠或优雅退出逻辑替代:
休眠方式(简单版)
#include <thread> #include <chrono> void Server::run(int port){ asio::ip::tcp::endpoint endpoint(asio::ip::tcp::v4(), port); std::cout << "[DEBUG] binded on " << port << std::endl; m_acceptor = std::make_shared<asio::ip::tcp::acceptor>(m_io, endpoint); listen(); io_thread = std::thread([&]{ m_io.run(); }); // 每秒休眠一次,减少CPU占用 while (true){ std::this_thread::sleep_for(std::chrono::seconds(1)); } }
优雅退出方式(推荐)
#include <csignal> #include <condition_variable> std::condition_variable cv; std::mutex mtx; bool running = true; void signal_handler(int signum){ running = false; cv.notify_one(); } void Server::run(int port){ // 注册Ctrl+C信号处理函数 std::signal(SIGINT, signal_handler); asio::ip::tcp::endpoint endpoint(asio::ip::tcp::v4(), port); std::cout << "[DEBUG] binded on " << port << std::endl; m_acceptor = std::make_shared<asio::ip::tcp::acceptor>(m_io, endpoint); listen(); io_thread = std::thread([&]{ m_io.run(); }); // 等待退出信号 std::unique_lock<std::mutex> lock(mtx); cv.wait(lock, []{ return !running; }); // 停止io_service并等待线程退出 m_io.stop(); if (io_thread.joinable()){ io_thread.join(); } }
4. 禁用Nagle算法(可选)
如果消息延迟是因为Nagle算法合并小数据包,可以在客户端连接时禁用该算法,减少发送延迟:
void Server::onAccept(std::shared_ptr<User> pUser){ // 禁用Nagle算法 asio::ip::tcp::no_delay option(true); pUser->socket.set_option(option); // 添加用户到m_users(需实现generate_user_id逻辑) int user_id = generate_user_id(); m_users[user_id] = pUser; // 发起异步读 asio::async_read(pUser->socket, asio::buffer(pUser->read_buf, PACK_SIZE), std::bind(&Server::onMessage, this, user_id, std::placeholders::_1)); // 继续监听新连接 listen(); }
内容的提问来源于stack exchange,提问作者kotivas
相关产品推荐
相关产品推荐

