Skip to content

为UDPRecv增加udp socket接收缓冲区大小调节功能 #296

Description

@dyx2025

为UDPRecv增加udp socket接收缓冲区大小调节功能

一个phxpaxos::Node只一个udp socket用于接收paxos udp消息,容易成为瓶颈。本人在开发机发现,连续发送多个udp消息且每个udp消息为64 000字节到同一个远端udp socket时,远端udp socket比较容易丢消息。调大远端udp socket接收缓冲区大小能有效缓解丢消息的情况。建议为UDPRecv增加udp socket接收缓冲区大小调节功能。

尽管调节UDPRecv udp socket接收缓冲区大小能有效缓解丢消息的情况,但并不代表业务层不会超时。因此业务侧需要观察调节后的耗时情况,如果不符合预期则回退调节操作。

修改后的代码路径:
include/phxpaxos/options.h

class Options
{ 
    // ...
    // 新增代码
    // 期望UDPRecv udp socket的接收缓冲区大小
    // < 0,不调节接收缓冲区大小
    // >= 0,表示用户期望设置的接收缓冲区大小,操作系统内核会根据其预设的下界和上界对接收缓冲区大小作进一步调整
    int  iUDPRecvUDPSocketRecvBuffSize;
};

phxpaxos/src/comm/options.cpp

Options :: Options()
{
    // ...
    iUDPRecvUDPSocketRecvBuffSize = -1;
};

src/node/pnode.cpp

int PNode :: InitNetWork(const Options & oOptions, NetWork *& poNetWork)
{                                                                                                                                                                                                           
    if (oOptions.poNetWork != nullptr)
    {   
        poNetWork = oOptions.poNetWork;
        PLImp("OK, use user network");
        return 0;
    }   

    // 修改代码
    /*
    int ret = m_oDefaultNetWork.Init(
            oOptions.oMyNode.GetIP(), oOptions.oMyNode.GetPort(), oOptions.iIOThreadCount);
    */
    int ret = m_oDefaultNetWork.Init(
            oOptions.oMyNode.GetIP(), oOptions.oMyNode.GetPort(), oOptions.iIOThreadCount, oOptions.iUDPRecvUDPSocketRecvBuffSize);
    if (ret != 0)
    {   
        PLErr("init default network fail, listenip %s listenport %d ret %d",
                oOptions.oMyNode.GetIP().c_str(), oOptions.oMyNode.GetPort(), ret);
        return ret;
    }   

    poNetWork = &m_oDefaultNetWork;
    
    PLImp("OK, use default network");

    return 0;
}

src/communicate/dfnetwork.h

class DFNetWork : public NetWork
{
    // ...
    // 修改代码
    // int Init(const std::string & sListenIp, const int iListenPort, const int iIOThreadCount);
    int Init(const std::string & sListenIp, const int iListenPort, const int iIOThreadCount, const int iUDPRecvUDPSocketRecvBuffSize);
    // ...
};

src/communicate/dfnetwork.cpp

// 修改代码
// int DFNetWork :: Init(const std::string & sListenIp, const int iListenPort, const int iIOThreadCount) 
int DFNetWork :: Init(const std::string & sListenIp, const int iListenPort, const int iIOThreadCount, const int iUDPRecvUDPSocketRecvBuffSize) 
{
    int ret = m_oUDPSend.Init();
    if (ret != 0)
    {
        return ret;
    }

    // 修改代码
    // ret = m_oUDPRecv.Init(iListenPort);
    ret = m_oUDPRecv.Init(iListenPort, iUDPRecvUDPSocketRecvBuffSize);
    if (ret != 0)
    {
        return ret;
    }

    ret = m_oTcpIOThread.Init(sListenIp, iListenPort, iIOThreadCount);
    if (ret != 0)
    {
        PLErr("m_oTcpIOThread Init fail, ret %d", ret);
        return ret;
    }

    return 0;
}

src/communicate/udp.h

class UDPRecv : public Thread
{
    // ...
    // 修改代码
    // int Init(const int iPort);
    int Init(const int iPort, const int iUDPRecvUDPSocketRecvBuffSize);
    // ...
};

src/communicate/udp.cpp

// 修改代码
// int UDPRecv :: Init(const int iPort)
int UDPRecv :: Init(const int iPort, const int iUDPRecvUDPSocketRecvBuffSize)
{
    if ((m_iSockFD = socket(AF_INET, SOCK_DGRAM, 0)) < 0)  
    {   
        return -1; 
    }   

    struct sockaddr_in addr;
    memset(&addr, 0, sizeof(addr));

    addr.sin_family = AF_INET;
    addr.sin_port = htons(iPort);
    addr.sin_addr.s_addr = htonl(INADDR_ANY);

    // 新增代码
    if (iUDPRecvUDPSocketRecvBuffSize >= 0) {
        int iUDPRecvUDPSocketRecvBuffSizeInput = iUDPRecvUDPSocketRecvBuffSize / 2;
        int ret = setsockopt(m_iSockFD , SOL_SOCKET, SO_RCVBUF, &iUDPRecvUDPSocketRecvBuffSizeInput , (socklen_t)sizeof(iUDPRecvUDPSocketRecvBuffSizeInput));
        if (-1 == ret) {
            PLErr("set receive buffer size fail");
            close(m_iSockFD);
            m_iSockFD = -1;
            return -1; 
        }
    }

    int enable = 1;
    setsockopt(m_iSockFD, SOL_SOCKET, SO_REUSEADDR, &enable, sizeof(int));

    if (bind(m_iSockFD, (struct sockaddr *)&addr, sizeof(addr)) < 0)
    {
        return -1;
    }

    return 0;
}

测试程序test_udp_sendto和test_udp_recvfrom分别为udp消息发送方和接收方。每次实验先重启test_udp_recvfrom,再重启test_udp_sendto,然后收集test_udp_recvfrom接收的udp消息数和test_udp_sendto发送的udp消息数。实验期间,每条udp消息固定为64000个字节。单次实验,test_udp_recvfrom使用udp socket默认接收缓冲区大小或者最大接收缓冲区大小,test_udp_sendto连续发送udp消息条数为200或2000或20000条。

本人开发机udp socket默认接收缓冲区大小为212992字节,udp socket最大接收缓冲区大小为8388608字节。

实验结果如下:
test_udp_recvfrom使用udp socket默认接收缓冲区大小,test_udp_sendto连续发送200条udp消息,丢消息率为0.3550-0.5800。
test_udp_recvfrom使用udp socket最大接收缓冲区大小,test_udp_sendto连续发送200条udp消息,丢消息率为0。

test_udp_recvfrom使用udp socket默认接收缓冲区大小,test_udp_sendto连续发送2000条udp消息,丢消息率为0.1075-0.2010。
test_udp_recvfrom使用udp socket最大接收缓冲区大小,test_udp_sendto连续发送2000条udp消息,丢消息率为0。

test_udp_recvfrom使用udp socket默认接收缓冲区大小,test_udp_sendto连续发送20000条udp消息,丢消息率为0.1318-0.1773。
test_udp_recvfrom使用udp socket最大接收缓冲区大小,test_udp_sendto连续发送20000条udp消息,丢消息率为0.0002-0.0027。

实验结果表明,调节udp socket接收缓冲区大小能有效缓解丢消息的情况。

测试代码:
guarder.h

#pragma once                                                                                                                                                                                                

#include <functional>

namespace test {

class Guarder {
public:
    // explicit Guarder(std::function<void()>&& func);
    Guarder(std::function<void()>&& func);
    ~Guarder();

    Guarder(const Guarder& other) = delete;
    Guarder& operator = (const Guarder& other) = delete;

    Guarder(Guarder&& other) = delete;
    Guarder& operator = (Guarder&& other) = delete;

private:
    std::function<void()> func_;
};

}

guarder.cpp

#include "guarder.h"                                                                                                                                                                                        

namespace test {

Guarder::Guarder(std::function<void()>&& func) {
    func_ = std::move(func);
}

Guarder::~Guarder() {
    func_();
}

}

test_udp_recvfrom.cpp

#include <sys/types.h> 
#include <sys/socket.h>
#include <strings.h>
#include <sys/time.h>
#include <netinet/in.h>
#include <arpa/inet.h>
#include <stdlib.h>
#include <errno.h>
#include <stdint.h>
#include <limits.h>
#include <string.h>
#include <assert.h>
#include <unistd.h>

#include <iostream>

#include "guarder.h"

int main(int argc, char** argv) {
    if (argc < 4) {
        std::cout << "argc < 4" << std::endl;
        return -1;
    }

    const char* ip = argv[1];
    int port = atoi(argv[2]);
    int use_max_recv_buff_size = atoi(argv[3]);

    std::cout << "ip = " << ip << " port = " << port << " use_max_recv_buff_size = " << use_max_recv_buff_size << std::endl; 

    int fd = socket(AF_INET, SOCK_DGRAM, 0);
    if (-1 == fd) {
        std::cout << "socket fail" << std::endl;
        return -1;
    }

    test::Guarder fd_guarder([&fd]() -> void {
        if (fd >= 0) {
            std::cout << "close fd = " << fd << std::endl;
            int ret = close(fd);
            assert(0 == ret);
            fd = -1;
        }
    });

    int ret = -1;

    int old_real_recv_buffer_size = 0;
    socklen_t addrlen = sizeof(old_real_recv_buffer_size);
    ret = getsockopt(fd, SOL_SOCKET, SO_RCVBUF, &old_real_recv_buffer_size, &addrlen);
    if (-1 == ret) {
        std::cout << "SO_RCVBUF: getsockopt fail" << std::endl;
        return -1;
    }

    std::cout << "old_real_recv_buffer_size = " << old_real_recv_buffer_size << std::endl;

    if (use_max_recv_buff_size) {
        int recv_buffer_size_input_paramter = INT_MAX / 2;
        ret = setsockopt(fd, SOL_SOCKET, SO_RCVBUF, &recv_buffer_size_input_paramter, (socklen_t)sizeof(recv_buffer_size_input_paramter));
        if (-1 == ret) {
            std::cout << "SO_RCVBUF: setsockopt fail" << std::endl;
            return -1;
        }

        int new_real_recv_buffer_size = 0;
        addrlen = sizeof(new_real_recv_buffer_size);
        ret = getsockopt(fd, SOL_SOCKET, SO_RCVBUF, &new_real_recv_buffer_size, &addrlen);
        if (-1 == ret) {
            std::cout << "SO_RCVBUF: getsockopt fail" << std::endl;
            return -1;
        }

        std::cout << "new_real_recv_buffer_size = " << new_real_recv_buffer_size << std::endl;
    }

    int reuse_addr = 1;
    ret = setsockopt(fd, SOL_SOCKET, SO_REUSEADDR, &reuse_addr, (socklen_t)sizeof(reuse_addr));
    if (-1 == ret) {
        std::cout << "SO_REUSEADDR: setsockopt fail" << std::endl;
        return -1;
    }

    sockaddr_in addr;
    bzero(&addr, sizeof(addr));
    addr.sin_family = AF_INET;
    addr.sin_addr.s_addr = inet_addr(ip);
    addr.sin_port = htons(port);

    ret = bind(fd, (const struct sockaddr*)&addr, (socklen_t)sizeof(addr));
    if (-1 == ret) {
        std::cout << "bind fail" << std::endl;
        return -1;
    }

    uint32_t succ_recv_count = 0;
    char buffer[64 * 1024];
    timeval start_tv;
    timeval end_tv;
    while (1) {
        ret = gettimeofday(&start_tv, NULL);
        if (-1 == ret) {
            std::cout << "start_tv: gettimeofday fail" << std::endl;
            return -1;
        }

        int recv_ret = recvfrom(fd, buffer, sizeof(buffer), 0, NULL, NULL);
        int recv_errno = errno;
        errno = 0;

        ret = gettimeofday(&end_tv, NULL);
        if (-1 == ret) {
            std::cout << "start_tv: gettimeofday fail" << std::endl;
            return -1;
        }

        std::cout << "recv_ret = " << recv_ret << " recv_errno = " << recv_errno 
            << " elapsed time us = " << end_tv.tv_sec * 1000000 + end_tv.tv_usec - start_tv.tv_sec * 1000000 - start_tv.tv_usec << std::endl;

        if (recv_ret != -1) {
            uint32_t net_send_number;
            assert(recv_ret >= sizeof(net_send_number));
            ++succ_recv_count;
            memcpy(&net_send_number, buffer, sizeof(net_send_number));
            uint32_t local_send_number = ntohl(net_send_number);
            std::cout << "succ_recv_count = " << succ_recv_count << " local_send_number = " << local_send_number << std::endl;
        }

    }

    std::cout << "finished" << std::endl;
    return 0;
}

test_udp_sendto.cpp

#include <sys/types.h> 
#include <sys/socket.h>
#include <strings.h>
#include <sys/time.h>
#include <netinet/in.h>
#include <arpa/inet.h>
#include <stdlib.h>
#include <errno.h>
#include <stdint.h>
#include <string.h>
#include <assert.h>
#include <unistd.h>

#include <iostream>

#include "guarder.h"

int main(int argc, char** argv) {
    if (argc < 3) {
        std::cout << "argc < 4" << std::endl;
        return -1;
    }

    const char* ip = argv[1];
    int port = atoi(argv[2]);
    uint32_t max_send_count = (uint32_t)atoi(argv[3]);

    std::cout << "ip = " << ip << " port = " << port << "max_send_count = " << max_send_count << std::endl; 

    int fd = socket(AF_INET, SOCK_DGRAM, 0);
    if (-1 == fd) {
        std::cout << "socket fail" << std::endl;
        return -1;
    }

    test::Guarder fd_guarder([&fd]() -> void {
        if (fd >= 0) {
            std::cout << "close fd = " << fd << std::endl;
            int ret = close(fd);
            assert(0 == ret);
            fd = -1;
        }
    });

    sockaddr_in addr;
    bzero(&addr, sizeof(addr));
    addr.sin_family = AF_INET;
    addr.sin_addr.s_addr = inet_addr(ip);
    addr.sin_port = htons(port);

    int ret = connect(fd, (const struct sockaddr*)&addr, (socklen_t)sizeof(addr));
    if (-1 == ret) {
        std::cout << "connect fail" << std::endl;
        return -1;
    }

    uint32_t succ_send_count = 0;
    uint32_t total_send_count = 0;
    char buffer[64 * 1000];
    while (total_send_count < max_send_count) {
        uint32_t net_send_number = htonl(total_send_count);
        memcpy(buffer, &net_send_number, sizeof(net_send_number));

        int send_ret = sendto(fd, buffer, sizeof(buffer), 0, (sockaddr*)&addr, (socklen_t)sizeof(addr));
        int send_errno = errno;

        std::cout << "send_ret = " << send_ret << " send_errno = " << send_errno << std::endl;

        if (send_ret != -1) {
            assert(send_ret >= sizeof(net_send_number));
            ++succ_send_count;
            std::cout << "succ_send_count = " << succ_send_count << std::endl;
        }

        ++total_send_count;
    }

    return 0;
}

编译命令:

g++ -O2 -std=c++11 -o test_udp_recvfrom guarder.cpp test_udp_recvfrom.cpp
g++ -O2 -std=c++11 -o test_udp_sendto guarder.cpp test_udp_sendto.cpp

运行命令:

# 运行test_udp_recvfrom
# 127.0.0.1表示ip,8888表示port,0表示使用udp socket默认接收缓冲区大小。把0改为1表示使用udp socket最大接收缓冲区大小。
./test_udp_recvfrom 127.0.0.1 8888 0

# 运行test_udp_sendto
# 127.0.0.1表示ip,8888表示port,200表示连续发送的udp消息条数。可以把200改为其他值,例如上述实验的2000或20000。
./test_udp_sendto 127.0.0.1 8888 200

Activity

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Metadata

Metadata

Assignees

No one assigned

    Labels

    No labels
    No labels

    Type

    No type

    Projects

    No projects

      Milestone

      No milestone

      Relationships

      None yet

      Development

      No branches or pull requests

      Issue actions