Boost Asio:如何在不同线程中实现异步TCP服务器的消息收发?
Boost.Asio异步TCP服务器线程通信解决方案
嘿,你这个基础框架已经搭得挺扎实了!针对你提到的两个线程通信问题,我给你梳理下具体的实现思路和代码修改方案:
1. 从主循环向TCP服务器线程发送数据
你尝试用asio::post()的思路完全正确!Boost.Asio的核心原则就是所有IO对象的操作必须在其所属的io_service线程上下文执行,所以用post()把写操作投递到io_service的线程里是标准做法。不过现在你的代码里有个问题:server类创建session后没有保存连接实例,导致没法主动给客户端发消息。
修改步骤:
- 给
server类添加一个线程安全的容器,用来管理所有活跃的session; - 给
server提供一个发送数据的接口,内部用asio::post()把写操作投递到io_service线程; 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来实现高效的等待,避免主线程空轮询。
实现步骤:
- 定义一个线程安全的消息队列类;
- 修改
session,收到数据后把消息放到队列; - 主线程启动一个循环,等待并处理队列中的消息。
首先实现线程安全队列:
#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
相关产品推荐
相关产品推荐

