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

Boost Asio:如何在不同线程中实现异步TCP服务器的消息收发?

Boost.Asio异步TCP服务器线程通信解决方案

嘿,你这个基础框架已经搭得挺扎实了!针对你提到的两个线程通信问题,我给你梳理下具体的实现思路和代码修改方案:


1. 从主循环向TCP服务器线程发送数据

你尝试用asio::post()的思路完全正确!Boost.Asio的核心原则就是所有IO对象的操作必须在其所属的io_service线程上下文执行,所以用post()把写操作投递到io_service的线程里是标准做法。不过现在你的代码里有个问题:server类创建session后没有保存连接实例,导致没法主动给客户端发消息。

修改步骤:

  1. 给server类添加一个线程安全的容器,用来管理所有活跃的session;
  2. 给server提供一个发送数据的接口,内部用asio::post()把写操作投递到io_service线程;
  3. session销毁时要从容器中移除自己,避免悬空引用。

修改后的关键代码片段:

首先修改server类:

#include <set>
#include <mutex>

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

  // 新增:向所有客户端发送数据的接口
  void send_to_all(const std::string& msg) {
    std::lock_guard<std::mutex> lock(sessions_mutex_);
    for (auto& session : sessions_) {
      // 用post把写操作投递到io_service线程
      boost::asio::post(session->get_io_context(), [session, msg]() {
        session->write(msg);
      });
    }
  }

  // 新增:添加session到管理列表
  void add_session(std::shared_ptr<session> s) {
    std::lock_guard<std::mutex> lock(sessions_mutex_);
    sessions_.insert(s);
  }

  // 新增:从管理列表移除session
  void remove_session(std::shared_ptr<session> s) {
    std::lock_guard<std::mutex> lock(sessions_mutex_);
    sessions_.erase(s);
  }

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

  tcp::acceptor acceptor_;
  tcp::socket socket_;
  std::set<std::shared_ptr<session>> sessions_;
  std::mutex sessions_mutex_;
};

然后修改session类,添加write方法和获取io_context的接口,以及在销毁时通知server移除自己:

class session : public std::enable_shared_from_this<session> {
public:
  session(tcp::socket socket, std::weak_ptr<server> srv) 
    : socket_(std::move(socket)), server_(srv) { }
  
  void start() { do_read(); }

  // 新增:主动写数据的方法
  void write(const std::string& msg) {
    auto self(shared_from_this());
    boost::asio::async_write(socket_, boost::asio::buffer(msg),
        [this, self](boost::system::error_code ec, std::size_t /*length*/) {
          if (ec) {
            // 写失败,移除session
            if (auto srv = server_.lock()) {
              srv->remove_session(self);
            }
          }
        });
  }

  boost::asio::io_context& get_io_context() {
    return socket_.get_executor().context();
  }

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);
          } else {
            // 读失败,移除session
            if (auto srv = server_.lock()) {
              srv->remove_session(self);
            }
          }
        });
  }

  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();
          } else {
            if (auto srv = server_.lock()) {
              srv->remove_session(self);
            }
          }
        });
  }

  tcp::socket socket_;
  std::weak_ptr<server> server_; // 用weak_ptr避免循环引用
  enum { max_length = 1024 };
  char data_[max_length];
};

然后在main函数里,你就可以这样主动发消息了:

// 比如主线程收到UI/命令行输入后,调用server的send_to_all
std::string input;
while (std::getline(std::cin, input)) {
  s->send_to_all(input);
}

2. 将TCP收到的数据传递到主函数处理

这里推荐用线程安全的消息队列来实现,核心思路是:

  • session收到数据后,把数据封装成消息,放到线程安全的队列里;
  • 主线程循环监听队列,有消息就拿出来处理;
  • 用std::condition_variable来实现高效的等待,避免主线程空轮询。

实现步骤:

  1. 定义一个线程安全的消息队列类;
  2. 修改session,收到数据后把消息放到队列;
  3. 主线程启动一个循环,等待并处理队列中的消息。

首先实现线程安全队列:

#include <queue>
#include <mutex>
#include <condition_variable>

template<typename T>
class thread_safe_queue {
public:
  void push(T msg) {
    std::lock_guard<std::mutex> lock(mtx_);
    queue_.push(std::move(msg));
    cv_.notify_one(); // 通知等待的线程
  }

  bool try_pop(T& out) {
    std::lock_guard<std::mutex> lock(mtx_);
    if (queue_.empty()) return false;
    out = std::move(queue_.front());
    queue_.pop();
    return true;
  }

  void wait_and_pop(T& out) {
    std::unique_lock<std::mutex> lock(mtx_);
    cv_.wait(lock, [this](){ return !queue_.empty(); });
    out = std::move(queue_.front());
    queue_.pop();
  }

  bool empty() {
    std::lock_guard<std::mutex> lock(mtx_);
    return queue_.empty();
  }

private:
  std::queue<T> queue_;
  std::mutex mtx_;
  std::condition_variable cv_;
};

然后修改session的构造函数和do_read方法,把收到的数据放到队列:

class session : public std::enable_shared_from_this<session> {
public:
  session(tcp::socket socket, std::weak_ptr<server> srv, thread_safe_queue<std::string>& msg_queue) 
    : socket_(std::move(socket)), server_(srv), msg_queue_(msg_queue) { }
  
  // ... 其他代码不变 ...

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) {
            // 把收到的数据放到消息队列
            std::string received(data_, length);
            msg_queue_.push(received);
            do_write(length);
          } else {
            if (auto srv = server_.lock()) {
              srv->remove_session(self);
            }
          }
        });
  }

  // ... 其他成员 ...
  thread_safe_queue<std::string>& msg_queue_;
};

同步修改server的do_accept方法,把消息队列传递给新创建的session:

// 先给server添加消息队列成员
class server : public std::enable_shared_from_this<server> {
public:
  server(boost::asio::io_service& io_service, short port, thread_safe_queue<std::string>& msg_queue)
    : acceptor_(io_service, tcp::endpoint(tcp::v4(), port)),
      socket_(io_service), msg_queue_(msg_queue) {
    do_accept();
  }

  // ... 其他代码不变 ...

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

  // ... 其他成员 ...
  thread_safe_queue<std::string>& msg_queue_;
};

最后在main函数里,主线程处理消息:

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;
    thread_safe_queue<std::string> msg_queue; // 创建消息队列
    std::shared_ptr<server> s = std::make_shared<server>(io_service, std::atoi(argv[1]), msg_queue);
    std::shared_ptr<boost::asio::io_service::work> work(new boost::asio::io_service::work(io_service));
    std::thread t1([&io_service]() {io_service.run();});
    
    // 主线程处理消息的循环
    std::string received_msg;
    while (true) {
      msg_queue.wait_and_pop(received_msg);
      std::cout << "Received from client: " << received_msg << std::endl;
      // 这里可以把消息传递给UI/命令行处理逻辑
    }

    t1.join();
  } catch (std::exception& e) {
    std::cerr << "Exception: " << e.what() << "\n";
  }
  return 0;
}

关键注意点

  • 避免线程直接操作IO对象:所有对socket的读写操作必须通过asio::post()投递到io_service线程执行,否则会有线程安全问题;
  • 用智能指针管理生命周期:std::shared_ptr和std::weak_ptr配合使用,避免循环引用和悬空指针;
  • 线程安全容器/队列:跨线程传递数据必须保证线程安全,std::mutex和std::condition_variable是标准的解决方案,比共享内存简单得多。

内容的提问来源于stack exchange,提问作者user_cr

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.06 19:12:48