Boost ASIO async_read未读取服务器响应问题求助
嘿,我来帮你梳理下这个问题的核心原因,然后给出具体的修复方案!你的场景里客户端能正常发数据给服务器,但服务器返回的"Response"客户端收不到,而且客户端的异步读取回调完全没触发,咱们一步步拆解:
核心问题拆解
客户端侧的关键问题
IO服务没有持续运行:
客户端构造函数里,connect()调用了m_ios.run(),但这时候还没有任何异步任务,run()会立即返回。之后你调用StartHandlingServer启动了async_read,但没有再次启动IO服务的事件循环——Boost.Asio的异步操作必须依赖io_service::run()(或run_one()/poll())来驱动回调执行,没跑事件循环的话,回调永远不会被触发。响应接收逻辑未实现:
receiveResponse()方法是空的,哪怕IO服务跑起来了,你也没写实际读取服务器响应的代码,自然拿不到数据。Socket初始化写法不规范:
socket_((new mysock(m_ios)))这种写法虽然能跑,但不符合shared_ptr的常规初始化方式,改成socket_(new mysock(m_ios))更清晰。
服务器侧的关键问题
Service对象生命周期失控:
在Accept()里你写了(new Service)->StartHandligClient(sock);——创建了Service对象但没保留任何指针,StartHandligClient返回后这个对象就会被销毁!而异步读取的lambda捕获了this,后续回调里访问this就是访问已经析构的对象,属于未定义行为,大概率会导致程序崩溃或逻辑异常。同步IO阻塞线程:
在异步读取header的回调里,你用了同步的asio::read读取数据部分,这会直接阻塞当前线程,违背了异步IO的设计初衷,还会导致其他异步任务无法及时处理。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

