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

Boost ASIO async_read未读取服务器响应问题求助

排查Boost.Asio客户端无法接收服务器响应的问题

嘿,我来帮你梳理下这个问题的核心原因,然后给出具体的修复方案!你的场景里客户端能正常发数据给服务器,但服务器返回的"Response"客户端收不到,而且客户端的异步读取回调完全没触发,咱们一步步拆解:

核心问题拆解

客户端侧的关键问题

  1. IO服务没有持续运行:
    客户端构造函数里,connect()调用了m_ios.run(),但这时候还没有任何异步任务,run()会立即返回。之后你调用StartHandlingServer启动了async_read,但没有再次启动IO服务的事件循环——Boost.Asio的异步操作必须依赖io_service::run()(或run_one()/poll())来驱动回调执行,没跑事件循环的话,回调永远不会被触发。

  2. 响应接收逻辑未实现:
    receiveResponse()方法是空的,哪怕IO服务跑起来了,你也没写实际读取服务器响应的代码,自然拿不到数据。

  3. Socket初始化写法不规范:
    socket_((new mysock(m_ios)))这种写法虽然能跑,但不符合shared_ptr的常规初始化方式,改成socket_(new mysock(m_ios))更清晰。

服务器侧的关键问题

  1. Service对象生命周期失控:
    在Accept()里你写了(new Service)->StartHandligClient(sock);——创建了Service对象但没保留任何指针,StartHandligClient返回后这个对象就会被销毁!而异步读取的lambda捕获了this,后续回调里访问this就是访问已经析构的对象,属于未定义行为,大概率会导致程序崩溃或逻辑异常。

  2. 同步IO阻塞线程:
    在异步读取header的回调里,你用了同步的asio::read读取数据部分,这会直接阻塞当前线程,违背了异步IO的设计初衷,还会导致其他异步任务无法及时处理。

  3. IO服务运行逻辑错误:
    Server::Run()里每次Accept()后调用m_ios.run(),但run()会在所有pending任务完成后立即返回,之后循环再去Accept(),但如果有新的异步任务(比如下一轮的read_async_1),IO服务没在运行,根本处理不了这些任务。

具体修复方案

客户端修复代码

我们调整为用线程持续运行IO服务,同时完善同步接收响应的逻辑(如果你想坚持异步读取,也可以基于这个基础修改):

#include <boost/asio.hpp>
#include <boost/archive/text_iarchive.hpp>
#include <boost/archive/text_oarchive.hpp>
#include <boost/bind.hpp>
#include <boost/serialization/vector.hpp>
#include <iostream>
#include <boost/shared_ptr.hpp>
#include <thread>
#include "../stocks.hpp"
using namespace boost;

class mysock : public boost::asio::ip::tcp::socket {
public:
 mysock(asio::io_service& serv) : boost::asio::ip::tcp::socket(serv) {}
 ~mysock() { std::cout << "Inside destructor for mysock \n"; }
};

class SyncTCPClient {
public:
 SyncTCPClient(const std::string& raw_ip_address, unsigned short port_num)
 : socket_(new mysock(m_ios)), // 修正shared_ptr初始化
 m_ep(asio::ip::address::from_string(raw_ip_address), port_num)
 {
 (*socket_.get()).open(m_ep.protocol());
 connect();
 // 启动线程持续运行IO服务
 m_io_thread.reset(new std::thread([this]() {
 m_ios.run();
 }));
 }

 ~SyncTCPClient() {
 // 停止IO服务并等待线程结束
 m_ios.stop();
 if (m_io_thread.joinable()) {
 m_io_thread.join();
 }
 }

 mysock& socket() { return *socket_.get(); }

 void connect() {
 (*socket_.get()).connect(m_ep);
 }

 void close() {
 (*socket_.get()).shutdown(boost::asio::ip::tcp::socket::shutdown_both);
 (*socket_.get()).close();
 }

 std::string emulateLongComputationOp(unsigned int duration_sec) {
 sendRequest("EMULATE_LONG_COMP_OP " + std::to_string(duration_sec) + "\n");
 return receiveResponse();
 };

private:
 void sendRequest(const std::string& /*request*/) {
 std::vector<stock> stocks_;
 stock s;
 s.code = "ABC";
 s.name = "A Big Company";
 s.open_price = 4.56;
 s.high_price = 5.12;
 s.low_price = 4.33;
 s.last_price = 4.98;
 s.buy_price = 4.96;
 s.buy_quantity = 1000;
 s.sell_price = 4.99;
 s.sell_quantity = 2000;
 stocks_.push_back(s);

 std::ostringstream archive_stream;
 boost::archive::text_oarchive archive(archive_stream);
 archive << stocks_;
 outbound_data_ = archive_stream.str();

 std::ostringstream header_stream;
 header_stream << std::setw(header_length) << std::hex << outbound_data_.size();
 if (!header_stream || header_stream.str().size() != header_length) {
 return;
 }
 outbound_header_ = header_stream.str();

 std::size_t headerSize = asio::write(*socket_.get(), boost::asio::buffer(outbound_header_));
 std::size_t dataSize = asio::write(*socket_.get(), boost::asio::buffer(outbound_data_));
 std::cout << "headerSize : " << headerSize << " , dataSize : " << dataSize << "\n";
 }

 std::string receiveResponse() {
 std::string response;
 asio::streambuf buf;
 // 同步读取服务器发送的"Response\n",直到换行符
 asio::read_until(*socket_.get(), buf, '\n');
 std::istream input(&buf);
 std::getline(input, response);
 return response;
 }

private:
 asio::io_service m_ios;
 std::unique_ptr<std::thread> m_io_thread; // 新增:运行IO服务的线程
 boost::shared_ptr<mysock> socket_;
 asio::ip::tcp::endpoint m_ep;
 enum { header_length = 8 };
 std::string outbound_data_;
 std::string outbound_header_;
};

int main() {
 const std::string raw_ip_address = "127.0.0.1";
 const unsigned short port_num = 3333;
 try {
 SyncTCPClient client(raw_ip_address, port_num);
 std::cout << "Sending request to the server... \n"<< std::endl;
 std::string response = client.emulateLongComputationOp(10);
 std::cout << "\nResponse received: " << response << std::endl;
 std::this_thread::sleep_for(std::chrono::seconds(10));
 std::cout << "\n\n Closing client connection \n\n";
 client.close();
 } catch (system::system_error &e) {
 std::cout << "Client Error occured! Error code = " << e.code() << ". Message: " << e.what();
 return e.code().value();
 }
 return 0;
}

服务器修复代码

我们用shared_ptr管理Service对象,把同步读取改成异步读取,同时修正IO服务的运行逻辑:

#include <boost/asio.hpp>
#include <boost/archive/text_iarchive.hpp>
#include <boost/archive/text_oarchive.hpp>
#include <boost/bind.hpp>
#include <boost/serialization/vector.hpp>
#include <boost/tuple/tuple.hpp>
#include <thread>
#include <atomic>
#include <memory>
#include <iostream>
#include "../stocks.hpp"
using namespace boost;

class Service {
public:
 Service(){} 
 void StartHandligClient(boost::shared_ptr<asio::ip::tcp::socket> sock) {
 std::cout << "StartHandligClient : sock.use_count : " << sock.use_count() << "\n";
 read_async_1(sock);
 }

private:
 void read_async_1(boost::shared_ptr<asio::ip::tcp::socket> sock) {
 if(!sock->is_open()) {
 std::cout << getpid() << " : Socket closed in sync_read \n" << std::flush;
 return ;
 }
 std::cout << "haha_1\n" << std::flush;
 boost::asio::async_read(
 *sock,
 boost::asio::buffer(inbound_header_),
 [this, sock](boost::system::error_code ec, size_t bytesRead) {
 int headerBytesReceived = bytesRead;
 std::cout << "\n\n headerBytesReceived : " << headerBytesReceived << "\n" << std::flush ;
 if (!ec) {
 std::istringstream is(std::string(inbound_header_, header_length));
 std::cout << "is : +" << is.str() << "+, inbound_header_ : +" << inbound_header_ << "\n";
 std::size_t inbound_data_size = 0;
 if (!(is >> std::hex >> inbound_data_size)) {
 std::cout << "RET-1 \n";
 return;
 }
 std::cout << "inbound_data_size : " << inbound_data_size << "\n" << std::flush;
 inbound_data_.resize(inbound_data_size);
 // 把同步read改成异步read
 boost::asio::async_read(
 *sock,
 boost::asio::buffer(inbound_data_),
 [this, sock](boost::system::error_code ec, size_t bytesRead) {
 if (!ec) {
 int bytesReceived = bytesRead;
 std::string archive_data(&inbound_data_[0], inbound_data_.size());
 std::istringstream archive_stream(archive_data);
 boost::archive::text_iarchive archive(archive_stream);
 archive >> stocks_;
 std::cout << "bytesReceived : " << bytesReceived << " , stocks_.size() : " << stocks_.size() << "\n";

 // 打印股票数据
 for (std::size_t i = 0; i < stocks_.size(); ++i) {
 std::cout << "Stock number " << i << "\n";
 std::cout << " code: " << stocks_[i].code << "\n";
 std::cout << " name: " << stocks_[i].name << "\n";
 std::cout << " open_price: " << stocks_[i].open_price << "\n";
 std::cout << " high_price: " << stocks_[i].high_price << "\n";
 std::cout << " low_price: " << stocks_[i].low_price << "\n";
 std::cout << " last_price: " << stocks_[i].last_price << "\n";
 std::cout << " buy_price: " << stocks_[i].buy_price << "\n";
 std::cout << " buy_quantity: " << stocks_[i].buy_quantity << "\n";
 std::cout << " sell_price: " << stocks_[i].sell_price << "\n";
 std::cout << " sell_quantity: " << stocks_[i].sell_quantity << "\n";
 }

 std::this_thread::sleep_for(std::chrono::seconds(1));
 // 异步发送响应
 std::string response = "Response\n";
 boost::asio::async_write(*sock, boost::asio::buffer(response), [this, sock](const boost::system::error_code& ec, size_t) {
 if (!ec) {
 this->read_async_1(sock); // 响应发送完成后继续等待下一个请求
 } else {
 std::cout << "Send response error: " << ec.message() << "\n";
 }
 });
 } else {
 std::cout << "Read data error: " << ec.message() << "\n";
 }
 }
 );
 } else {
 if(ec == boost::asio::error::eof) {
 std::cout << getpid() << " : ** sync_read : Connection lost : boost::asio::error::eof ** \n";
 }
 std::cout << "Error occured in async_read! Error code = " << ec.value() << ". Message: " << ec.message() << "\n" << std::flush;
 }
 }
 );
 }

 enum { header_length = 8 };
 char inbound_header_[header_length];
 std::vector<char> inbound_data_;
 std::vector<stock> stocks_;
};

class Acceptor {
public:
 Acceptor(asio::io_service& ios, unsigned short port_num)
 : m_ios(ios),
 m_acceptor(m_ios, asio::ip::tcp::endpoint(asio::ip::address_v4::any(), port_num))
 {
 m_acceptor.listen();
 start_accept(); // 启动异步Accept
 }

private:
 void start_accept() {
 boost::shared_ptr<asio::ip::tcp::socket> sock(new asio::ip::tcp::socket(m_ios));
 m_acceptor.async_accept(*sock, [this, sock](const boost::system::error_code& ec) {
 if (!ec) {
 // 用shared_ptr管理Service,确保对象生命周期足够长
 boost::shared_ptr<Service> service(new Service);
 service->StartHandligClient(sock);
 }
 start_accept(); // 继续等待下一个连接
 });
 }

 asio::io_service& m_ios;
 asio::ip::tcp::acceptor m_acceptor;
};

class Server {
public:
 Server() : m_stop(false) {}
 void Start(unsigned short port_num) {
 m_thread.reset(new std::thread([this, port_num]() {
 Run(port_num);
 }));
 }
 void Stop() {
 std::cout << "STOPPING \n";
 m_stop.store(true);
 m_ios.stop(); // 停止IO服务
 m_thread->join();
 }

private:
 void Run(unsigned short port_num) {
 Acceptor acc(m_ios, port_num);
 m_ios.run(); // 持续运行IO服务,处理所有异步任务
 }

 std::unique_ptr<std::thread> m_thread;
 std::atomic<bool> m_stop;
 asio::io_service m_ios;
};

int main() {
 unsigned short port_num = 3333;
 try {
 Server srv;
 srv.Start(port_num);
 std::this_thread::sleep_for(std::chrono::seconds(100));
 std::cout << "Stopping server \n";
 srv.Stop();
 } catch (system::system_error &e) {
 std::cout << "Error occured! Error code = " << e.code() << ". Message: " << e.what();
 }
 return 0;
}

修复后效果

  • 客户端启动后会持续运行IO服务,发送请求后能同步读取到服务器返回的"Response"。
  • 服务器用异步Accept和异步IO处理请求,Service对象由shared_ptr管理,不会提前析构,响应能正常发送给客户端。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.15 03:47:03