Boost.Asio TCP服务端高并发连接报错:remote_endpoint未连接
使用Boost.Asio实现TCP服务端与客户端通信程序,单个或少量客户端连接时服务端接收数据正常,但每秒接入100个及以上客户端时,会出现“remote_endpoint: Transport endpoint is not connected”错误,寻求解决方案。
#include <iostream> #include <array> #include <map> #include <boost/asio.hpp> #include <mysql/mysql.h> #include <nlohmann/json.hpp> #include <chrono> #include <cmath> using namespace boost::asio; using json = nlohmann::json; class AsyncTCPServer { public: // Constructor AsyncTCPServer(io_service& io_service, short port) : acceptor_(io_service, ip::tcp::endpoint(ip::tcp::v4(), port)), socket_(io_service) { StartAccept(); ConnectToMariaDB(); } // Destructor ~AsyncTCPServer() { // Disconnect from MariaDB if (db_connection_ != nullptr) { mysql_close(db_connection_); std::cout << "MariaDB connection closed" << std::endl; } } private: // Start accepting new client connections and initiate asynchronous communication void StartAccept() { acceptor_.async_accept(socket_, [this](boost::system::error_code ec) { if(ec) { std::cerr << "Error occurred during connection acceptance: " << ec.message() << std::endl; } else { // Add the new client socket to the array for management AddClient(std::make_shared<ip::tcp::socket>(std::move(socket_))); // Wait for the next client connection StartAccept(); // Initiate asynchronous communication for the current client StartRead(clients_[num_clients_ - 1]); } }); } // Add a new client to the array and print a connection message void AddClient(std::shared_ptr<ip::tcp::socket> client) { std::lock_guard<std::mutex> lock(mutex_); if (num_clients_ < max_clients) { clients_[num_clients_] = client; num_clients_++; // Print connection message std::cout << client->remote_endpoint().address().to_string() + " connected." << std::endl; } else { std::cerr << "Cannot accept connection, maximum number of clients exceeded." << std::endl; client->close(); } } // Start asynchronous reading for the client void StartRead(std::shared_ptr<ip::tcp::socket> client) { auto& buffer = buffers_[client]; // Get the buffer associated with the client async_read_until(*client, buffer, '\n', [this, client, &buffer](boost::system::error_code ec, std::size_t length) { if (ec) { RemoveClient(client); } else { std::istream is(&buffer); std::string message; std::getline(is, message); // Save to the database SaveToDatabase(message); StartRead(client); } }); } // Connect to MariaDB void ConnectToMariaDB() { // Initialize MariaDB connection db_connection_ = mysql_init(nullptr); if (db_connection_ == nullptr) { std::cerr << "Failed to initialize MariaDB connection" << std::endl; exit(1); } const char* host = "localhost"; const char* user = "root"; const char* password = "1234"; const char* database = "servertest"; unsigned int port = 3306; if (mysql_real_connect(db_connection_, host, user, password, database, port, nullptr, 0) == nullptr) { std::cerr << "Failed to connect to MariaDB: " << mysql_error(db_connection_) << std::endl; exit(1); } std::cout << "MariaDB connection established" << std::endl; } // Save message to the database void SaveToDatabase(const std::string& message) { // Assume JSON and save to the database try { auto j = json::parse(message); std::lock_guard<std::mutex> lock(mutex_); if (j.find("CPU") != j.end()) { SaveCpuToDatabase(j["CPU"]); } if (j.find("NIC") != j.end()) { SaveNicToDatabase(j["NIC"]); } if (j.find("Memory") != j.end()) { SaveMemoryToDatabase(j["Memory"]); } if (j.find("Disk") != j.end()) { SaveDiskToDatabase(j["Disk"]); } std::cout << "Saved to the database" << std::endl; } catch (const nlohmann::detail::parse_error& e) { // Catch JSON parsing errors std::cerr << "JSON parsing error: " << e.what() << std::endl; std::cerr << "Error occurred in the message: " << message << std::endl; } catch (const std::exception& e) { std::cerr << "Failed to save to the database: " << e.what() << std::endl; std::cerr << "Error occurred in the message: " << message << std::endl; } } void SaveCpuToDatabase(const json& cpuData) { for (const auto& processor : cpuData) { // Extract information for each processor std::string cores = processor["Cores"].get<std::string>(); std::string model = processor["Model"].get<std::string>(); std::string siblings = processor["Siblings"].get<std::string>(); // Generate query and save data to the DB std::string query = "INSERT INTO cpu_table (cores, model, siblings) VALUES ('" + cores + "', '" + model + "', '" + siblings + "')"; if (mysql_query(db_connection_, query.c_str()) != 0) { std::cerr << "Error occurred while saving CPU information to the database: " << mysql_error(db_connection_) << std::endl; } } } void SaveNicToDatabase(const json& nicData) { for (const auto& nic : nicData) { // Extract information for each NIC std::string interface = nic["Interface"].get<std::string>(); std::string mac_address = nic["MAC Address"].get<std::string>(); std::string operational_state = nic["Operational State"].get<std::string>(); std::string speed = nic["Speed"].get<std::string>(); // Generate query and save data to the DB std::string query = "INSERT INTO nic_table (interface, mac_address, operational_state, speed) VALUES ('" + interface + "', '" + mac_address + "', '" + operational_state + "', '" + speed + "')"; int queryResult = mysql_query(db_connection_, query.c_str()); if (queryResult != 0) { std::cerr << "Error occurred while saving NIC information to the database: " << mysql_error(db_connection_) << std::endl; } } } void SaveMemoryToDatabase(const json& memoryData) { // Similar logic for saving memory information to the database // ... } void SaveDiskToDatabase(const json& diskData) { // Similar logic for saving disk information to the database // ... } // Remove client from the array and print a connection termination message void RemoveClient(std::shared_ptr<ip::tcp::socket> client) { for (int i = 0; i < max_clients; ++i) { if (clients_[i] == client) { clients_[i] = nullptr; num_clients_--; // Print connection termination message std::cout << client->remote_endpoint().address().to_string() + " connection terminated." << std::endl; boost::system::error_code ec; client->shutdown(ip::tcp::socket::shutdown_both, ec); client->close(ec); break; } } } ip::tcp::acceptor acceptor_; // TCP acceptor ip::tcp::socket socket_; // TCP socket static const int max_clients = 100000; // Maximum number of clients std::array<std::shared_ptr<ip::tcp::socket>, max_clients> clients_; // TCP sockets corresponding to the maximum number of clients int num_clients_ = 0; // Current number of connected clients std::mutex mutex_; // MariaDB connection handler MYSQL* db_connection_; std::map<std::shared_ptr<ip::tcp::socket>, streambuf> buffers_; // Map to manage buffers for each client }; int main() { std::chrono::system_clock::time_point start = std::chrono::system_clock::now(); try { boost::asio::io_service io_service; // Create io_service object AsyncTCPServer server(io_service, 12345); // Create an object of the AsyncTCPServer class (object, port number) io_service.run(); // Start the event loop } catch (std::exception& e) { std::cerr << "Exception caught: " << e.what() << std::endl; std::chrono::duration<double> sec = std::chrono::system_clock::now() - start; std::cout << "Time taken to run (seconds): " << sec.count() << " seconds" << std::endl; } return 0; }
1. 共享socket引发的竞态条件
问题根源
代码使用成员变量socket_作为所有async_accept操作的目标socket,高并发下多个异步accept操作复用同一个socket对象,前一个连接未处理完成时,后一个accept会覆盖socket状态,导致socket连接无效,调用remote_endpoint触发错误。
解决方法
每次async_accept创建独立socket对象,避免共享:
void StartAccept() { auto new_socket = std::make_shared<ip::tcp::socket>(acceptor_.get_executor().context()); acceptor_.async_accept(*new_socket, [this, new_socket](boost::system::error_code ec) { if(ec) { std::cerr << "Error occurred during connection acceptance: " << ec.message() << std::endl; } else { AddClient(new_socket); StartAccept(); StartRead(new_socket); } }); }
删除类中成员变量ip::tcp::socket socket_;。
2. 断开连接后调用remote_endpoint
问题根源
RemoveClient中直接调用client->remote_endpoint(),此时客户端可能已断开连接,socket处于未连接状态,触发错误。
解决方法
连接成功时提前保存远程地址,避免断开后再获取:
// 类中新增存储客户端地址的容器 std::map<std::shared_ptr<ip::tcp::socket>, std::string> client_addresses_; void AddClient(std::shared_ptr<ip::tcp::socket> client) { std::lock_guard<std::mutex> lock(mutex_); boost::system::error_code ec; auto ep = client->remote_endpoint(ec); if (ec) { std::cerr << "Failed to get remote endpoint: " << ec.message() << std::endl; client->close(); return; } std::string addr = ep.address().to_string(); if (num_clients_ < max_clients) { clients_[num_clients_] = client; client_addresses_[client] = addr; num_clients_++; std::cout << addr << " connected." << std::endl; } else { std::cerr << "Cannot accept connection, maximum number of clients exceeded." << std::endl; client->close(); } } void RemoveClient(std::shared_ptr<ip::tcp::socket> client) { std::lock_guard<std::mutex> lock(mutex_); std::string addr = client_addresses_[client]; client_addresses_.erase(client); for (int i = 0; i < max_clients; ++i) { if (clients_[i] == client) { clients_[i] = nullptr; num_clients_--; std::cout << addr << " connection terminated." << std::endl; boost::system::error_code ec; client->shutdown(ip::tcp::socket::shutdown_both, ec); client->close(ec); break; } } }
3. 同步数据库操作阻塞IO线程
问题根源
SaveToDatabase及数据库插入操作同步执行,阻塞Boost.Asio的IO线程,高并发下IO线程无法及时处理新连接和读写请求,导致客户端连接超时断开,引发socket状态异常。
解决方法
将数据库操作转移到独立线程池异步执行:
// 类中新增线程池 boost::asio::thread_pool db_pool(4); // 修改SaveToDatabase void SaveToDatabase(const std::string& message) { boost::asio::post(db_pool, [this, message]() { try { auto j = json::parse(message); std::lock_guard<std::mutex> lock(mutex_); if (j.find("CPU") != j.end()) { SaveCpuToDatabase(j["CPU"]); } if (j.find("NIC") != j.end()) { SaveNicToDatabase(j["NIC"]); } if (j.find("Memory") != j.end()) { SaveMemoryToDatabase(j["Memory"]); } if (j.find("Disk") != j.end()) { SaveDiskToDatabase(j["Disk"]); } std::cout << "Saved to the database" << std::endl; } catch (const nlohmann::detail::parse_error& e) { std::cerr << "JSON parsing error: " << e.what() << std::endl; std::cerr << "Error occurred in the message: " << message << std::endl; } catch (const std::exception& e) { std::cerr << "Failed to save to the database: " << e.what() << std::endl; std::cerr << "Error occurred in the message: " << message << std::endl; } }); } // 析构函数等待线程池完成任务 ~AsyncTCPServer() { db_pool.join(); if (db_connection_ != nullptr) { mysql_close(db_connection_); std::cout << "MariaDB connection closed" << std::endl; } }
4. 客户端管理容器效率低下
问题根源
使用std::array存储客户端socket,RemoveClient时需遍历整个数组,高并发下锁竞争加剧,降低处理效率。
解决方法
替换为std::unordered_set,实现O(1)插入和删除:
// 替换原有的clients_和num_clients_ std::unordered_set<std::shared_ptr<ip::tcp::socket>> clients_; void AddClient(std::shared_ptr<ip::tcp::socket> client) { std::lock_guard<std::mutex> lock(mutex_); boost::system::error_code ec; auto ep = client->remote_endpoint(ec); if (ec) { std::cerr << "Failed to get remote endpoint: " << ec.message() << std::endl; client->close(); return; } std::string addr = ep.address().to_string(); if (clients_.size() < max_clients) { clients_.insert(client); client_addresses_[client] = addr; std::cout << addr << " connected." << std::endl; } else { std::cerr << "Cannot accept connection, maximum number of clients exceeded." << std::endl; client->close(); } } void RemoveClient(std::shared_ptr<ip::tcp::socket> client) { std::lock_guard<std::mutex> lock(mutex_); std::string addr = client_addresses_[client]; client_addresses_.erase(client); clients_.erase(client); std::cout << addr << " connection terminated." << std::endl; boost::system::error_code ec; client->shutdown(ip::tcp::socket::shutdown_both, ec); client->close(ec); }
5. 客户端buffer的线程安全问题
问题根源
多线程运行io_service时,多个线程可能同时访问同一客户端的streambuf,引发竞态条件。
解决方法
使用boost::asio::strand序列化每个客户端的读写操作:
// 定义客户端数据结构体 struct ClientData { std::shared_ptr<ip::tcp::socket> socket; boost::asio::strand<boost::asio::io_context::executor

