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(); }); }); }
修复方案
问题根源
- Session的成员变量(
message_size_、buffer_)未同步,多个异步操作并发访问时引发数据竞争,导致未定义行为(如挂起)。 - 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
相关产品推荐
相关产品推荐

