【问题标题】:Boost asio TCP async server not async?Boost asio TCP异步服务器不是异步的?
【发布时间】:2014-12-22 17:06:06
【问题描述】:

我正在使用Boost example中提供的代码。

服务器一次只接受 1 个连接。这意味着,在当前连接关闭之前,不会有新连接。

如何让上面的代码同时接受无限连接?

#include <cstdlib>
#include <iostream>
#include <memory>
#include <utility>
#include <boost/asio.hpp>

using boost::asio::ip::tcp;

class session
  : public std::enable_shared_from_this<session>
{
public:
  session(tcp::socket socket)
    : socket_(std::move(socket))
  {
  }

  void start()
  {
    do_read();
  }

private:
  void do_read()
  {
    auto self(shared_from_this());
    socket_.async_read_some(boost::asio::buffer(data_, max_length),
        [this, self](boost::system::error_code ec, std::size_t length)
        {
          if (!ec)
          {
            boost::this_thread::sleep(boost::posix_time::milliseconds(10000));//sleep some time
            do_write(length);
          }
        });
  }

  void do_write(std::size_t length)
  {
    auto self(shared_from_this());
    boost::asio::async_write(socket_, boost::asio::buffer(data_, length),
        [this, self](boost::system::error_code ec, std::size_t /*length*/)
        {
          if (!ec)
          {
            do_read();
          }
        });
  }

  tcp::socket socket_;
  enum { max_length = 1024 };
  char data_[max_length];
};

class server
{
public:
  server(boost::asio::io_service& io_service, short port)
    : acceptor_(io_service, tcp::endpoint(tcp::v4(), port)),
      socket_(io_service)
  {
    do_accept();
  }

private:
  void do_accept()
  {
    acceptor_.async_accept(socket_,
        [this](boost::system::error_code ec)
        {
          if (!ec)
          {
            std::make_shared<session>(std::move(socket_))->start();
          }

          do_accept();
        });
  }

  tcp::acceptor acceptor_;
  tcp::socket socket_;
};

int main(int argc, char* argv[])
{
  try
  {
    if (argc != 2)
    {
      std::cerr << "Usage: async_tcp_echo_server <port>\n";
      return 1;
    }

    boost::asio::io_service io_service;

    server s(io_service, std::atoi(argv[1]));

    io_service.run();
  }
  catch (std::exception& e)
  {
    std::cerr << "Exception: " << e.what() << "\n";
  }

  return 0;
}

如您所见,程序等待睡眠,在此期间它没有获取第二个连接。

【问题讨论】:

  • 如果您使用示例中的代码,那么它将接受多个连接。 async_accept 的 AcceptHandler 启动另一个 async_accept 操作。
  • 您如何确定它一次只接受 1 个连接?你测试了吗?如果是这样,请解释您是如何测试的以及您得到了什么结果。你是通过代码检查确定的吗?如果是,请解释你是如何得出这个结论的。
  • @Tanner Sansbury:我正在使用示例中的代码和一个非常简单的睡眠来测试它是否接受多个连接。
  • @David Schwartz:很抱歉没有早点发布我的代码,在这里。
  • @SpeedCoder 该“睡眠”不会测试它是否接受多个连接。如果它处于睡眠状态,它如何接受另一个连接?

标签: c++ c++11 boost asynchronous


【解决方案1】:

您正在处理程序内部进行同步等待,该处理程序在为您的 io_service 提供服务的唯一线程上运行。这使得 Asio 等待为任何新请求调用处理程序。

  1. 使用deadline_timewait_async,或者,

    void do_read() {
        auto self(shared_from_this());
        socket_.async_read_some(boost::asio::buffer(data_, max_length),
                                [this, self](boost::system::error_code ec, std::size_t length) {
            if (!ec) {
                timer_.expires_from_now(boost::posix_time::seconds(1));
                timer_.async_wait([this, self, length](boost::system::error_code ec) {
                        if (!ec)
                            do_write(length);
                    });
            }
        });
    }
    

    其中timer_ 字段是sessionboost::asio::deadline_timer 成员

  2. 作为穷人的解决方案,添加更多线程(这只是意味着如果同时到达的请求多于处理它们的线程,它仍然会阻塞,直到第一个线程变得可用于接收新请求)

    boost::thread_group tg;
    for (int i=0; i < 10; ++i)
        tg.create_thread([&]{ io_service.run(); });
    
    tg.join_all();
    

【讨论】:

  • 您能否解释一下如何在 boost 示例中初始化截止时间计时器?这对我来说似乎是不可能的......
  • 我了解多线程解决方案的工作原理,但我不明白如何使用截止时间计时器来实现。在类服务器中,我添加了一个名为:boost::asio::io_service *io_service;(这是指向 io_service 的指针)的成员,然后在调用 timer_.expires... 之前启动了计时器:deadline_timer timer_(*io_service); 这个编译器正常,但我总是得到错误:The I/O operation has been aborted because of either a thread exit or an application request. 任何想法请问为什么?
  • @SpeedCoder 看看我一个小时前链接的代码。它是完整的并且有效。可能你忘了打expires_from_now
【解决方案2】:

原始代码和修改后的代码都是异步的,并且接受多个连接。从下面的sn-p中可以看出,async_accept操作的AcceptHandler发起了另一个async_accept操作,形成了一个异步循环:

        .-----------------------------------.
        V                                   |
void server::do_accept()                    |
{                                           |
  acceptor_.async_accept(...,               |
      [this](boost::system::error_code ec)  |
      {                                     |
        // ...                              |
        do_accept();  ----------------------'
      });
}

session 的 ReadHandler 中的sleep() 导致运行io_service 的线程阻塞,直到睡眠完成。因此,该程序将什么也不做。但是,这不会导致取消任何未完成的操作。为了更好地理解异步操作和io_service,请考虑阅读this 答案。


这是一个示例demonstrating 处理多个连接的服务器。它产生一个线程,创建 5 个客户端套接字并将它们连接到服务器。

#include <cstdlib>
#include <iostream>
#include <memory>
#include <utility>
#include <vector>
#include <boost/asio.hpp>
#include <boost/thread.hpp>

using boost::asio::ip::tcp;

class session
  : public std::enable_shared_from_this<session>
{
public:
  session(tcp::socket socket)
    : socket_(std::move(socket))
  {
  }

  ~session()
  {
    std::cout << "session ended" << std::endl;
  }

  void start()
  {
    std::cout << "session started" << std::endl;
    do_read();
  }

private:
  void do_read()
  {
    auto self(shared_from_this());
    socket_.async_read_some(boost::asio::buffer(data_, max_length),
        [this, self](boost::system::error_code ec, std::size_t length)
        {
          if (!ec)
          {
            do_write(length);
          }
        });
  }

  void do_write(std::size_t length)
  {
    auto self(shared_from_this());
    boost::asio::async_write(socket_, boost::asio::buffer(data_, length),
        [this, self](boost::system::error_code ec, std::size_t /*length*/)
        {
          if (!ec)
          {
            do_read();
          }
        });
  }

  tcp::socket socket_;
  enum { max_length = 1024 };
  char data_[max_length];
};

class server
{
public:
  server(boost::asio::io_service& io_service, short port)
    : acceptor_(io_service, tcp::endpoint(tcp::v4(), port)),
      socket_(io_service)
  {
    do_accept();
  }

private:
  void do_accept()
  {
    acceptor_.async_accept(socket_,
        [this](boost::system::error_code ec)
        {
          if (!ec)
          {
            std::make_shared<session>(std::move(socket_))->start();
          }

          do_accept();
        });
  }

  tcp::acceptor acceptor_;
  tcp::socket socket_;
};

int main(int argc, char* argv[])
{
  try
  {
    if (argc != 2)
    {
      std::cerr << "Usage: async_tcp_echo_server <port>\n";
      return 1;
    }

    boost::asio::io_service io_service;

    auto port = std::atoi(argv[1]);
    server s(io_service, port);

    boost::thread client_main(
        [&io_service, port]
        {
          tcp::endpoint server_endpoint(
              boost::asio::ip::address_v4::loopback(), port);

          // Create and connect 5 clients to the server.
          std::vector<std::shared_ptr<tcp::socket>> clients;
          for (auto i = 0; i < 5; ++i)
          {
              auto client = std::make_shared<tcp::socket>(
                  std::ref(io_service));
              client->connect(server_endpoint);
              clients.push_back(client);
          }

          // Wait 2 seconds before destroying all clients.
          boost::this_thread::sleep(boost::posix_time::seconds(2));
        });

   io_service.run();
   client_main.join();
  }
  catch (std::exception& e)
  {
    std::cerr << "Exception: " << e.what() << "\n";
  }

  return 0;
}

输出:

session started
session started
session started
session started
session started
session ended
session ended
session ended
session ended
session ended

【讨论】:

    猜你喜欢
    • 1970-01-01
    • 2012-08-27
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2015-08-07
    相关资源
    最近更新 更多