Boost Asio多线程UDP客户端实现效率与技术问题咨询
using RawDataArray = std::array<unsigned char, 65000>; class StaticBuffer { private: RawDataArray m_data; std::size_t m_n_avail; public: StaticBuffer() : m_data(), m_n_avail(0) {} StaticBuffer(std::size_t n_bytes) { m_n_avail = n_bytes; } StaticBuffer(const StaticBuffer& other) { std::cout << "ctor cpy\n"; m_data = other.m_data; m_n_avail = other.m_n_avail; } StaticBuffer(const StaticBuffer& other, std::size_t n_bytes) { std::cout << "ctor cpy\n"; m_data = other.m_data; m_n_avail = n_bytes; } StaticBuffer(const RawDataArray& data, std::size_t n_bytes) { std::cout << "ctor static buff\n"; m_data = data; m_n_avail = n_bytes; } void set_size(int n) { m_n_avail = n; } void set_max_size() { m_n_avail = m_data.size(); } std::size_t max_size()const { return m_data.size(); } unsigned char& operator[](std::size_t i) { return m_data[i]; } const unsigned char& operator[] (std::size_t i)const { return m_data[i]; } StaticBuffer& operator=(const StaticBuffer& other) { if (this == &other) return *this; m_data = other.m_data; m_n_avail = other.m_n_avail; return *this; } void push_back(unsigned char val) { if (m_n_avail < m_data.size()) { m_data[m_n_avail] = val; } else throw "Out of memory"; } void reset() { m_n_avail = 0; } unsigned char* data() { return m_data.data(); } const unsigned char* data()const { return m_data.data(); } std::size_t size()const { return m_n_avail; } ~StaticBuffer() {} }; class UDPSeassion; using DataBuffer = StaticBuffer; using DataBufferPtr = std::unique_ptr<DataBuffer>; using ExternalReadHandler = std::function<void(DataBufferPtr)>; class UDPSeassion : public std::enable_shared_from_this<UDPSeassion> { private: asio::io_context& m_ctx; asio::ip::udp::socket m_socket; asio::ip::udp::endpoint m_endpoint; std::string m_addr; unsigned short m_port; asio::io_context::strand m_send_strand; std::deque<DataBufferPtr> m_dq_send; asio::io_context::strand m_rcv_strand; DataBufferPtr m_rcv_data; ExternalReadHandler external_rcv_handler; private: void do_send_data_from_dq() { if (m_dq_send.empty()) return; m_socket.async_send_to( asio::buffer(m_dq_send.front()->data(), m_dq_send.front()->size()), m_endpoint, asio::bind_executor(m_send_strand, [this](const boost::system::error_code& er, std::size_t bytes_transferred) { if (!er) { m_dq_send.pop_front(); do_send_data_from_dq(); } else { //post to loggger } })); } void do_read(const boost::system::error_code& er, std::size_t bytes_transferred) { if (!er) { m_rcv_data->set_size(bytes_transferred); asio::post(m_ctx, [this, data = std::move(m_rcv_data)]() mutable { external_rcv_handler(std::move(data)); }); m_rcv_data = std::make_unique<DataBuffer>(); m_rcv_data->set_max_size(); async_read(); } } public: UDPSeassion(asio::io_context& ctx, const std::string& addr, unsigned short port) : m_ctx(ctx), m_socket(ctx), m_endpoint(asio::ip::address::from_string(addr), port), m_addr(addr), m_port(port), m_send_strand(ctx), m_dq_send(), m_rcv_strand(ctx), m_rcv_data(std::make_unique<DataBuffer>(65000)) {} ~UDPSeassion() {} const std::string& get_host()const { return m_addr; } unsigned short get_port() { return m_port; } template<typename ExternalReadHandlerCallable> void set_read_data_headnler(ExternalReadHandlerCallable&& handler) { external_rcv_handler = std::forward<ExternalReadHandlerCallable>(handler); } void start() { m_socket.open(asio::ip::udp::v4()); async_read(); } void async_read() { m_socket.async_receive_from( asio::buffer(m_rcv_data->data(), m_rcv_data->size()), m_endpoint, asio::bind_executor(m_rcv_strand, std::bind(&UDPSeassion::do_read, this, std::placeholders::_1, std::placeholders::_2)) ); } void async_send(DataBufferPtr pData) { asio::post(m_ctx, asio::bind_executor(m_send_strand, [this, pDt = std::move(pData)]() mutable { m_dq_send.emplace_back(std::move(pDt)); if (m_dq_send.size() == 1) do_send_data_from_dq(); })); } }; void handler_read(DataBufferPtr pdata) { // decoding raw_data -> decod_data // lock mutext // queue.push_back(decod_data) // unlock mutext //for view pdata std::stringstream ss; ss << "thread handler: " << std::this_thread::get_id() << " " << pdata->data() << " " << pdata->size() << std::endl; std::cout << ss.str() << std::endl; } int main() { asio::io_context ctx; //auto work_guard = asio::make_work_guard(ctx); std::cout << "MAIN thread: " << std::this_thread::get_id() << std::endl; StaticBuffer b{4}; b[0]='A'; b[1]='B'; b[2]='C'; b[4]='\n'; UDPSeassion client(ctx,"127.0.0.1",11223); client.set_read_data_headnler(handler_read); client.start(); std::vector<std::thread> threads; for (int i=0;i<3;++i) { threads.emplace_back([&](){ std::stringstream ss; ss << "run thread: " << std::this_thread::get_id() << std::endl; std::cout << ss.str(); ctx.run(); std::cout << "end thread\n"; } ); } client.async_send(std::make_unique<StaticBuffer>(b)); ctx.run(); for (auto& t:threads) t.join(); return 1; }
技术问题
- 若将该类应用于运行频率约100Hz(每10ms发送一次系统状态)的系统中:
1.1 当前多线程实现是否合理?效率表现如何?
1.2 仅使用单线程处理读写的客户端实现效率怎样? - 通过
std::move(unique_ptr_data)在任务间传递缓冲区的方式是否正确? - 实际应用中,应为UDP客户端分配多少线程处理读写操作?
- 针对TCP客户端,上述线程分配问题的答案是什么?
问题解答
1. 100Hz系统下的实现分析
1.1 当前多线程实现的合理性与效率
当前实现用strand保护发送队列和接收缓冲区的访问,多线程运行io_context的设计是合理的,但对100Hz的场景来说属于性能过剩。
100Hz意味着每秒仅100次发送操作,单线程就能轻松处理。不过当前实现的优势在于:如果后续业务逻辑(比如handler_read里的解码、队列操作)比较耗时,多线程可以把这些阻塞/CPU密集的任务分摊到不同线程,避免阻塞IO事件循环。
效率方面,strand保证同一逻辑链的操作串行执行,无数据竞争;多线程处理不同任务时能利用多核资源。但如果业务逻辑本身很轻量,多线程带来的上下文切换开销反而会抵消优势。
1.2 单线程处理读写的效率
单线程实现对100Hz的场景完全够用,甚至更高效。因为没有多线程的上下文切换开销,io_context的事件循环能以极低延迟处理IO事件。
只要业务逻辑(比如解码、数据处理)不是长时间阻塞或CPU密集型,单线程的响应速度和资源占用都会更优。就算偶尔有短时间的计算任务,10ms的间隔也足够单线程完成,不会影响下一次发送。
2. std::move(unique_ptr_data)传递缓冲区的正确性
这种方式正确且推荐:
unique_ptr用于独占所有权,std::move能高效转移缓冲区的所有权,避免不必要的内存拷贝(你的StaticBuffer是65KB的数组,拷贝成本很高)。- 异步任务间传递时,所有权明确,不会出现悬空指针或重复释放的问题。比如
do_read里把m_rcv_datamove给post的任务,之后立刻创建新的接收缓冲区,逻辑清晰安全。
3. UDP客户端的线程分配建议
UDP是无连接协议,线程分配核心取决于业务逻辑的复杂度:
- 如果业务逻辑(读写后的处理)很轻量:分配1个线程完全足够,甚至不需要多线程。
- 如果业务逻辑有CPU密集型操作(比如大数据解码、复杂计算),或者有阻塞操作(比如磁盘IO、数据库访问):可以分配2~4个线程(和CPU核心数匹配即可,不用过多),让
io_context的事件循环和业务处理能并行执行。 - 极端场景下,如果有大量并发的UDP会话(比如同时和几十个设备通信),可以根据会话数量适当增加线程,但一般来说,1~4个线程就足够应对绝大多数UDP客户端场景。
4. TCP客户端的线程分配建议
TCP是面向连接的协议,线程分配逻辑和UDP类似,但有几点差异:
- 单线程依然适用于轻量业务的单个TCP连接,ASIO的异步模型能高效处理单连接的读写。
- 如果是多个TCP客户端连接,或者单个连接的业务逻辑有阻塞/CPU密集操作:建议线程数等于CPU核心数(比如4核就开4个线程),让
io_context能并行处理多个连接的IO事件,同时分摊业务处理的负载。 - 注意:TCP的每个连接的读写操作依然要用
strand保护(如果多线程运行io_context),避免同一连接的读写操作并发执行导致数据混乱。
内容的提问来源于stack exchange,提问作者Lucifer
相关产品推荐
相关产品推荐

