You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

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

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.06.28 14:24:53