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

C++ Boost Asio单客户端多请求并发处理问题求助

问题描述

基于C++ Boost开发的服务器对每个客户端只能串行处理请求,无法同时处理同一客户端的多个请求,仅能依次执行。尝试将do_read方法放在post之后,期望提交请求到线程池后立即启动下一次读取,但服务器接收任意一个请求后就挂起,无法响应其他请求。

相关代码如下:

Session类的do_read与read_message方法

void Session::do_read()
{
    auto self(shared_from_this());

    boost::asio::async_read(socket_, boost::asio::buffer(&message_size_, sizeof(message_size_)),
        [this, self](boost::system::error_code ec, std::size_t /*length*/)
        {
            if (!ec)
            {
                message_size_ = ntohl(message_size_);
                buffer_.resize(message_size_);
                read_message();
            }
            else
            {
                BOOST_LOG_TRIVIAL(error) << "Error reading message length: " << ec.message();
            }
        });
}

void Session::read_message()
{
    auto self(shared_from_this());

    boost::asio::async_read(socket_, boost::asio::buffer(buffer_),
                            [this, self](boost::system::error_code ec, std::size_t /*length*/)
                            {
                                if (!ec)
                                {
                                    std::string request_(buffer_.begin(), buffer_.end());
                                    nlohmann::json jsonRequest_ = jsonWorker_.parceJson(request_);

                                    boost::asio::post(threadPool_, [this, self, jsonRequest_](){
                                        requestRouter_.defineQuery(self->socket_.get_executor(), userID_, jsonRequest_, connectionPool_, shared_from_this(), connectedUsers_);
                                    });

                                    do_read();
                                }
                                else
                                {
                                    BOOST_LOG_TRIVIAL(error) << "Error reading message: " << ec.message() << "\n";
                                }
                            });
}

线程池执行的defineQuery方法中写socket的逻辑

if(json_["Info"] == "GET_ID")
{
    std::string responseJson_ = jsonWorker_.createUserIdJson(userID_);

    boost::asio::post(executor_, [this, session_, responseJson_](){
        session_->do_write(responseJson_);
    });
}

此前错误写法(仅能接收单个请求)

if(json_["Info"] == "GET_ID")
{
    std::string responseJson_ = jsonWorker_.createUserIdJson(userID_);

    boost::asio::post(executor_, [this, session_, responseJson_](){
        session_->do_read();
        session_->do_write(responseJson_);
    });
}

尝试移动do_read位置但仍挂起的写法

if (!ec)
{
    std::string request_(buffer_.begin(), buffer_.end());
    nlohmann::json jsonRequest_ = jsonWorker_.parceJson(request_);

    boost::asio::post(threadPool_, [this, self, jsonRequest_](){
        requestRouter_.defineQuery(self->socket_.get_executor(), userID_, jsonRequest_, connectionPool_, shared_from_this(), connectedUsers_);
    });

    boost::asio::post(self->socket_.get_executor(), [this, self](){
        do_read();
    });
}

if (!ec)
{
    std::string request_(buffer_.begin(), buffer_.end());
    nlohmann::json jsonRequest_ = jsonWorker_.parceJson(request_);

    boost::asio::post(threadPool_, [this, self, jsonRequest_](){
        requestRouter_.defineQuery(self->socket_.get_executor(), userID_, jsonRequest_, connectionPool_, shared_from_this(), connectedUsers_);

        boost::asio::post(self->socket_.get_executor(), [this, self](){
            do_read();
        });
    });
}
修复方案

问题根源

  1. Session的成员变量(message_size_、buffer_)未同步,多个异步操作并发访问时引发数据竞争,导致未定义行为(如挂起)。
  2. TCP是字节流协议,同一连接的读取必须串行执行,但请求的处理逻辑可以并行,之前的错误写法混淆了读取串行与处理并行的关系。

具体修复步骤

1. 用Strand保护Session的异步操作

给Session添加strand,确保所有socket操作的回调串行执行,避免成员变量的并发修改:

class Session : public std::enable_shared_from_this<Session> {
private:
    boost::asio::strand<boost::asio::io_context::executor_type> strand_;
    boost::asio::ip::tcp::socket socket_;
    uint32_t message_size_ = 0;
    std::vector<uint8_t> buffer_;
    // ... 其他成员变量

public:
    Session(boost::asio::ip::tcp::socket socket, /* 其他构造参数 */)
        : socket_(std::move(socket)),
          strand_(socket_.get_executor())
    {}

    // ... 其他方法
};

2. 修改do_read和read_message,绑定回调到Strand

所有异步操作的回调都通过boost::asio::bind_executor绑定到strand,确保串行执行:

void Session::do_read()
{
    auto self(shared_from_this());
    boost::asio::post(strand_, [this, self]() {
        boost::asio::async_read(socket_, boost::asio::buffer(&message_size_, sizeof(message_size_)),
            boost::asio::bind_executor(strand_,
                [this, self](boost::system::error_code ec, std::size_t /*length*/)
                {
                    if (!ec)
                    {
                        message_size_ = ntohl(message_size_);
                        buffer_.resize(message_size_);
                        read_message();
                    }
                    else
                    {
                        BOOST_LOG_TRIVIAL(error) << "Error reading message length: " << ec.message();
                    }
                }));
    });
}

void Session::read_message()
{
    auto self(shared_from_this());
    boost::asio::async_read(socket_, boost::asio::buffer(buffer_),
        boost::asio::bind_executor(strand_,
            [this, self](boost::system::error_code ec, std::size_t /*length*/)
            {
                if (!ec)
                {
                    // 复制请求数据,避免后续读取修改buffer影响处理逻辑
                    std::string request(buffer_.begin(), buffer_.end());
                    nlohmann::json jsonRequest = jsonWorker_.parceJson(request);

                    // 提交处理任务到线程池,捕获复制后的变量而非this,避免生命周期问题
                    boost::asio::post(threadPool_, [jsonRequest, self, userID = userID_, conn_pool = connectionPool_, connected_users = connectedUsers_](){
                        requestRouter_.defineQuery(self->strand_, userID, jsonRequest, conn_pool, self, connected_users);
                    });

                    // 立即启动下一次读取,strand保证串行安全
                    do_read();
                }
                else
                {
                    BOOST_LOG_TRIVIAL(error) << "Error reading message: " << ec.message();
                }
            }));
}

3. 修正do_write方法,同样用Strand保护

确保写操作也在strand上串行执行,避免并发写冲突:

void Session::do_write(const std::string& response)
{
    auto self(shared_from_this());
    boost::asio::post(strand_, [this, self, response]() {
        // 按协议打包响应:先写长度,再写内容
        uint32_t response_size = htonl(response.size());
        std::vector<uint8_t> write_buf;
        write_buf.reserve(sizeof(response_size) + response.size());
        
        // 写入长度
        std::copy(reinterpret_cast<const uint8_t*>(&response_size),
                  reinterpret_cast<const uint8_t*>(&response_size) + sizeof(response_size),
                  std::back_inserter(write_buf));
        // 写入响应内容
        std::copy(response.begin(), response.end(), std::back_inserter(write_buf));

        boost::asio::async_write(socket_, boost::asio::buffer(write_buf),
            boost::asio::bind_executor(strand_,
                [this, self](boost::system::error_code ec, std::size_t /*length*/)
                {
                    if (ec)
                    {
                        BOOST_LOG_TRIVIAL(error) << "Error writing message: " << ec.message();
                    }
                }));
    });
}

4. 调整defineQuery中的回调逻辑

确保提交到executor的写操作使用Session的strand:

if(json_["Info"] == "GET_ID")
{
    std::string responseJson_ = jsonWorker_.createUserIdJson(userID_);

    boost::asio::post(strand_, [session_, responseJson_](){
        session_->do_write(responseJson_);
    });
}

修复后效果

  • 同一客户端的请求读取过程串行(符合TCP流特性,不会出现数据错乱)。
  • 请求的处理逻辑在线程池中并行执行,实现同一客户端多个请求的同时处理。
  • Strand保护所有成员变量的访问,避免数据竞争,解决服务器挂起问题。

内容的提问来源于stack exchange,提问作者zxctatar

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.13 23:29:58