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

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
相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.05 06:23:26