【问题标题】:Boost asio and array of sockets提升 asio 和套接字阵列
【发布时间】:2015-04-27 06:14:01
【问题描述】:

我必须编写一个初始化 TCP 套接字数组的程序,并使用异步 i/o 使用线程池读取数据。我是异步 io、线程池、shared_ptrs 的新手。我现在拥有的是一个带有 one 套接字的工作程序。剪辑如下:

boost::shared_ptr< asio::ip::tcp::socket > sock1(
    new asio::ip::tcp::socket( *io_service )
);

    boost::shared_ptr< asio::ip::tcp::acceptor > acceptor( new asio::ip::tcp::acceptor( *io_service ) );
    asio::ip::tcp::endpoint endpoint(asio::ip::tcp::v4(), portNum);
    acceptor->open( endpoint.protocol() );
    acceptor->set_option( asio::ip::tcp::acceptor::reuse_address( false ) );
    acceptor->bind( endpoint );
    acceptor->listen();

我一直在为“套接字数组”获取类似的代码,也就是说,我想要绑定到端点[] 的接受器[]。我必须将指针传递给接受器和套接字,所以 shared_ptr 进来了,但我无法正确处理。

for (i=0; i<10; i++) {
    // init socket[i] with *io_service
    // init endpoint[i]
    // init acceptor[i] with *io_service
    acceptor[i]->listen()
    }

(顺便说一句,我真的需要一个 socket[] 数组来处理这个 porpose 吗?)有人可以帮我吗?

【问题讨论】:

    标签: c++ arrays sockets boost shared-ptr


    【解决方案1】:

    这是使用 Boost ASIO 实现 TCP 回显服务器侦听多个端口的完整示例,并使用线程池将工作分配到多个内核。它基于 Boost 文档中的 this 示例(提供单线程 TCP 回显服务器)。

    会话类

    session 类表示与客户端的单个活动套接字连接。它从套接字读取,然后将相同的数据写入套接字以将其回显给客户端。该实现使用 Boost ASIO 提供的 async_... 函数:这些函数在 I/O 服务处注册一个回调,该回调将在 I/O 操作完成时触发。

    session.h

    #pragma once
    
    #include <array>
    #include <memory>
    
    #include <boost/asio.hpp>
    
    /**
     * A TCP session opened on the server.
     */
    class session : public std::enable_shared_from_this<session> {
    
      using endpoint_t = boost::asio::ip::tcp::endpoint;
      using socket_t = boost::asio::ip::tcp::socket;
    
    public:
      session(boost::asio::io_service &service);
    
      /**
       * Start reading from the socket.
       */
      void start();
    
      /**
       * Callback for socket reads.
       */
      void handle_read(const boost::system::error_code &ec,
                       size_t bytes_transferred);
    
      /**
       * Callback for socket writes.
       */
      void handle_write(const boost::system::error_code &ec);
    
      /**
       * Get a reference to the session socket.
       */
      socket_t &socket() { return socket_; }
    
    private:
      /**
       * Session socket
       */
      socket_t socket_;
    
      /**
       * Buffer to be used for r/w operations.
       */
      std::array<uint8_t, 4096> buffer_;
    };
    

    session.cpp

    #include "session.h"
    
    #include <functional>
    #include <iostream>
    #include <thread>
    
    using boost::asio::async_write;
    using boost::asio::buffer;
    using boost::asio::io_service;
    using boost::asio::error::connection_reset;
    using boost::asio::error::eof;
    using boost::system::error_code;
    using boost::system::system_error;
    
    using std::placeholders::_1;
    using std::placeholders::_2;
    
    session::session(io_service &service) : socket_{service} {}
    
    void session::start() {
      auto handler = std::bind(&session::handle_read, shared_from_this(), _1, _2);
      socket_.async_read_some(buffer(buffer_), handler);
    }
    
    void session::handle_read(const error_code &ec, size_t bytes_transferred) {
      if (ec) {
        if (ec == eof || ec == connection_reset) {
          return;
        }
    
        throw system_error{ec};
      }
    
      std::cout << "Thread " << std::this_thread::get_id() << ": Received "
                << bytes_transferred << " bytes on " << socket_.local_endpoint()
                << " from " << socket_.remote_endpoint() << std::endl;
    
      auto handler = std::bind(&session::handle_write, shared_from_this(), _1);
      async_write(socket_, buffer(buffer_.data(), bytes_transferred), handler);
    }
    
    void session::handle_write(const error_code &ec) {
      if (ec) {
        throw system_error{ec};
      }
    
      auto handler = std::bind(&session::handle_read, shared_from_this(), _1, _2);
      socket_.async_read_some(buffer(buffer_), handler);
    }
    

    服务器类

    服务器类为每个给定端口创建一个接受器。接受器将侦听端口并为每个传入的连接请求分派一个套接字。再次使用 async_... 函数实现等待传入连接。

    server.h

    #pragma once
    
    #include <vector>
    
    #include <boost/asio.hpp>
    
    #include "session.h"
    
    /**
     * Listens to a socket and dispatches sessions for each incoming request.
     */
    class server {
    
      using acceptor_t = boost::asio::ip::tcp::acceptor;
      using endpoint_t = boost::asio::ip::tcp::endpoint;
      using socket_t = boost::asio::ip::tcp::socket;
    
    public:
      server(boost::asio::io_service &service, const std::vector<uint16_t> &ports);
    
      /**
       * Start listening for incoming requests.
       */
      void start_accept(size_t index);
    
      /**
       * Callback for when a request comes in.
       */
      void handle_accept(size_t index, std::shared_ptr<session> new_session,
                         const boost::system::error_code &ec);
    
    private:
    
      /**
       * Reference to the I/O service that will call our callbacks.
       */
      boost::asio::io_service &service_;
    
      /**
       * List of acceptors each listening to (a different) socket.
       */
      std::vector<acceptor_t> acceptors_;
    };
    

    server.cpp

    #include "server.h"
    
    #include <functional>
    
    #include <boost/asio.hpp>
    
    using std::placeholders::_1;
    using std::placeholders::_2;
    
    using boost::asio::io_service;
    using boost::asio::error::eof;
    using boost::system::error_code;
    using boost::system::system_error;
    
    server::server(boost::asio::io_service &service,
                   const std::vector<uint16_t> &ports)
        : service_{service} {
    
      auto create_acceptor = [&](uint16_t port) {
        acceptor_t acceptor{service};
        endpoint_t endpoint{boost::asio::ip::tcp::v4(), port};
        acceptor.open(endpoint.protocol());
        acceptor.set_option(acceptor_t::reuse_address(false));
        acceptor.bind(endpoint);
        acceptor.listen();
        return acceptor;
      };
    
      std::transform(ports.begin(), ports.end(), std::back_inserter(acceptors_),
                     create_acceptor);
    
      for (size_t i = 0; i < acceptors_.size(); i++) {
        start_accept(i);
      }
    }
    
    void server::start_accept(size_t index) {
      auto new_session{std::make_shared<session>(service_)};
    
      auto handler =
          std::bind(&server::handle_accept, this, index, new_session, _1);
    
      acceptors_[index].async_accept(new_session->socket(), handler);
    }
    
    void server::handle_accept(size_t index, std::shared_ptr<session> new_session,
                               const boost::system::error_code &ec) {
      if (ec) {
        throw system_error{ec};
      }
    
      new_session->start();
      start_accept(index);
    }
    

    主要

    主函数为一系列端口创建服务器。 对于此示例,端口设置为 5000,...,5010。然后,它为每个调用 Boost ASIO 提供的 I/O 服务的 run 函数的 CPU 内核生成一系列线程。 I/O 服务能够处理这样的多线程场景,在调用其run 函数(reference)的线程之间分派工作:

    多个线程可以调用 run() 函数来建立一个线程池,io_context 可以从中执行处理程序。在池中等待的所有线程都是等效的,io_context 可以选择其中任何一个来调用处理程序。

    server_main.cpp

    
    #include "server.h"
    
    #include <numeric>
    
    #include <boost/asio.hpp>
    #include <boost/thread.hpp>
    
    int main() {
      std::vector<uint16_t> ports{};
    
      // Fill ports with range [5000,5000+n)
      ports.resize(10);
      std::iota(ports.begin(), ports.end(), 5000);
    
      boost::asio::io_service service{};
    
      server s{service, ports};
    
      // Spawn thread group for running the I/O service
      size_t thread_count = std::min(
          static_cast<size_t>(boost::thread::hardware_concurrency()), ports.size());
    
      boost::thread_group tg{};
      for (size_t i = 0; i < thread_count; ++i) {
        tg.create_thread([&]() { service.run(); });
      }
    
      tg.join_all();
    
      return 0;
    }
    

    例如,您可以使用g++ -O2 -lboost_thread -lpthread {session,server,server_main}.cpp -o server 编译服务器。如果您使用向其发送随机数据的客户端运行服务器,您将获得如下输出:

    Thread 140043413878528: Received 4096 bytes on 127.0.0.1:5007 from 127.0.0.1:40856
    Thread 140043405485824: Received 4096 bytes on 127.0.0.1:5000 from 127.0.0.1:42556
    Thread 140043388700416: Received 4096 bytes on 127.0.0.1:5005 from 127.0.0.1:58582
    Thread 140043388700416: Received 4096 bytes on 127.0.0.1:5001 from 127.0.0.1:40192
    Thread 140043388700416: Received 4096 bytes on 127.0.0.1:5003 from 127.0.0.1:42508
    Thread 140043397093120: Received 4096 bytes on 127.0.0.1:5008 from 127.0.0.1:37808
    Thread 140043388700416: Received 4096 bytes on 127.0.0.1:5006 from 127.0.0.1:35440
    Thread 140043397093120: Received 4096 bytes on 127.0.0.1:5009 from 127.0.0.1:58306
    Thread 140043405485824: Received 4096 bytes on 127.0.0.1:5002 from 127.0.0.1:56300
    

    您可以看到服务器处理多个端口,工作分布在工作线程之间(不一定将每个线程限制到特定端口)。

    【讨论】:

      猜你喜欢
      • 1970-01-01
      • 1970-01-01
      • 2012-03-04
      • 1970-01-01
      • 2017-07-20
      • 1970-01-01
      • 2013-11-07
      • 1970-01-01
      • 2010-10-06
      相关资源
      最近更新 更多