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

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;
}
技术问题
  1. 若将该类应用于运行频率约100Hz(每10ms发送一次系统状态)的系统中:
    1.1 当前多线程实现是否合理?效率表现如何?
    1.2 仅使用单线程处理读写的客户端实现效率怎样?
  2. 通过std::move(unique_ptr_data)在任务间传递缓冲区的方式是否正确?
  3. 实际应用中,应为UDP客户端分配多少线程处理读写操作?
  4. 针对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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.14 03:08:10