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

如何解决UDP非阻塞socket调用sendto()时的EAGAIN错误

问题分析与解决方案

问题根源

非阻塞UDP socket出现EAGAIN错误,本质是发送缓冲区已满:单个设备发送时,数据量小,缓冲区能及时被内核排空;同时给4个设备连续发送时,短时间内的数据包总量超过了缓冲区的处理能力,内核无法立即接受新的发送请求,直接返回EAGAIN。

以下是针对性的解决方案,兼顾非阻塞要求、发送频率需求和设备公平性:


方案1:增大UDP发送缓冲区

先尝试提升socket的发送缓冲区上限,减少缓冲区满的概率。系统默认UDP发送缓冲区可能较小(通常在几KB级别),连续发送4个300字节的包(含UDP头部)容易触顶。

修改open_chamber_socket函数,添加SO_SNDBUF配置:

int open_chamber_socket(void)
{
    int code, flags;
    struct sockaddr_in addr;
    int optval = 1;
    
    e_socket = socket (AF_INET, SOCK_DGRAM, IPPROTO_UDP);
    if (e_socket < 0) {
        log_perror (_("%s: Can't create chamber socket."), MYNAME);
        return e_socket;
    }
    // 设置非阻塞标志
    flags = fcntl (e_socket, F_GETFL, 0);
    if (flags != -1)
        fcntl (e_socket, F_SETFL, flags | O_NONBLOCK);
        
    // 增大发送缓冲区到16KB(可根据实际调整)
    int send_buf_size = 16384;
    setsockopt(e_socket, SOL_SOCKET, SO_SNDBUF, &send_buf_size, sizeof(send_buf_size));
    
    setsockopt(e_socket, SOL_SOCKET, SO_REUSEADDR, &optval, sizeof(optval));
    bzero((char *) &addr, sizeof addr);
    addr.sin_family = AF_INET;
    addr.sin_addr.s_addr = htonl (CHAMBER_ADDRESS);
    addr.sin_port = htons (CHAMBER_PORT);
    code = bind (e_socket, (struct sockaddr *)&addr, sizeof addr);
    if (code < 0) {
        log_perror (_("%s: Can't bind chamber socket to 10.0.0.1:%d."), MYNAME, CHAMBER_PORT);
        return code;
    }
    return 0;
}

注意:系统有缓冲区最大限制,可通过sysctl net.core.wmem_max查看,若设置值超过该限制,内核会自动降到最大值。


方案2:用select()等待socket可写,轮询发送保证公平性

不能阻塞整个程序,但可以在发送前用select()检查socket是否可写,避免盲目调用sendto触发EAGAIN。同时采用轮询方式发送,确保每个设备都有平等的发送机会。

修改deferred_activities的发送逻辑:

deferred_activities() {
    ... 
    // 省略构建发送缓冲区的代码
    ...
    int j;
    fd_set write_fds;
    struct timeval tv;

    for (j = 0; j < 4; j++) {
        FD_ZERO(&write_fds);
        FD_SET(e_socket, &write_fds);
        // 设置10ms超时,避免长时间阻塞影响其他功能
        tv.tv_sec = 0;
        tv.tv_usec = 10000;

        int select_ret = select(e_socket + 1, NULL, &write_fds, NULL, &tv);
        if (select_ret == -1) {
            if (errno != EINTR) {
                log_perror(MYNAME ": select error when waiting for socket write");
            }
            continue;
        } else if (select_ret == 0) {
            // 超时跳过当前设备,保证其他设备能被处理
            continue;
        }

        // socket可写,执行发送
        struct sockaddr_in addr = {.sin_family = AF_INET};
        addr.sin_addr.s_addr = device[j]->net_address;
        addr.sin_port = device[j]->net_port;
        int code = sendto(e_socket, device[j]->buf, device[j]->datasize, 0, (struct sockaddr *) &addr, sizeof addr);
        if (code == -1) {
            if (errno != EAGAIN) {
                log_perror(MYNAME ": sendto error for device %d", j);
            }
            // 若仍遇EAGAIN,本次跳过,下次重试
        }
    }
}

方案3:分散发送频率,避免集中发送

当前deferred_activities每次调用都发送4个设备的包,容易瞬间填满缓冲区。可以改为每次只发送一个设备,轮询循环,分散发送压力,同时控制发送频率满足每秒2-3条的要求。

// 全局静态变量,记录当前发送的设备索引
static int current_device_idx = 0;

deferred_activities() {
    ... 
    // 省略构建发送缓冲区的代码
    ...
    struct sockaddr_in addr = {.sin_family = AF_INET};
    addr.sin_addr.s_addr = device[current_device_idx]->net_address;
    addr.sin_port = device[current_device_idx]->net_port;
    
    fd_set write_fds;
    struct timeval tv;
    FD_ZERO(&write_fds);
    FD_SET(e_socket, &write_fds);
    // 设置5ms超时
    tv.tv_sec = 0;
    tv.tv_usec = 5000;

    int select_ret = select(e_socket + 1, NULL, &write_fds, NULL, &tv);
    if (select_ret > 0 && FD_ISSET(e_socket, &write_fds)) {
        int code = sendto(e_socket, device[current_device_idx]->buf, device[current_device_idx]->datasize, 0, (struct sockaddr *) &addr, sizeof addr);
        if (code == -1 && errno != EAGAIN) {
            log_perror(MYNAME ": sendto error for device %d", current_device_idx);
        }
    }

    // 切换到下一个设备,循环轮询
    current_device_idx = (current_device_idx + 1) % 4;
}

频率调整:当前主循环每500us轮询一次,每12次调用deferred_activities(即每6ms调用一次)。若要每个设备每秒发送2-3条,可将模数改为40(每20ms调用一次),4次调用覆盖所有设备,刚好满足每秒2.5次的发送频率。


方案4:添加发送队列,重试EAGAIN的数据包

对于确实遇到EAGAIN的数据包,将其加入队列,后续调用时重试,避免消息丢失。

// 定义发送队列结构体
typedef struct {
    struct sockaddr_in addr;
    unsigned char *buf;
    size_t datasize;
    int retry_count; // 最大重试次数
} SendQueueItem;

#define MAX_QUEUE_SIZE 16
static SendQueueItem send_queue[MAX_QUEUE_SIZE];
static int queue_head = 0;
static int queue_tail = 0;

// 添加数据包到发送队列
int add_to_send_queue(const struct sockaddr_in *addr, const unsigned char *buf, size_t datasize) {
    int next_tail = (queue_tail + 1) % MAX_QUEUE_SIZE;
    if (next_tail == queue_head) {
        log_warning(MYNAME ": send queue full, dropping packet");
        return -1;
    }
    send_queue[queue_tail].addr = *addr;
    send_queue[queue_tail].buf = malloc(datasize);
    if (!send_queue[queue_tail].buf) {
        log_perror(MYNAME ": malloc failed for send queue");
        return -1;
    }
    memcpy(send_queue[queue_tail].buf, buf, datasize);
    send_queue[queue_tail].datasize = datasize;
    send_queue[queue_tail].retry_count = 3;
    queue_tail = next_tail;
    return 0;
}

// 处理发送队列中的待发数据包
void process_send_queue() {
    fd_set write_fds;
    struct timeval tv;
    FD_ZERO(&write_fds);
    FD_SET(e_socket, &write_fds);
    tv.tv_sec = 0;
    tv.tv_usec = 2000;

    int select_ret = select(e_socket + 1, NULL, &write_fds, NULL, &tv);
    if (select_ret <= 0) {
        return;
    }

    while (queue_head != queue_tail && FD_ISSET(e_socket, &write_fds)) {
        SendQueueItem *item = &send_queue[queue_head];
        int code = sendto(e_socket, item->buf, item->datasize, 0, (struct sockaddr *)&item->addr, sizeof(item->addr));
        if (code == -1) {
            if (errno == EAGAIN) {
                // 仍不可写,退出循环下次再试
                break;
            } else {
                log_perror(MYNAME ": sendto failed in queue");
                free(item->buf);
                queue_head = (queue_head + 1) % MAX_QUEUE_SIZE;
            }
        } else {
            // 发送成功,移除队列
            free(item->buf);
            queue_head = (queue_head + 1) % MAX_QUEUE_SIZE;
        }

        // 再次检查socket状态
        FD_ZERO(&write_fds);
        FD_SET(e_socket, &write_fds);
        select_ret = select(e_socket + 1, NULL, &write_fds, NULL, &tv);
        if (select_ret <= 0) {
            break;
        }
    }
}

// 修改deferred_activities
deferred_activities() {
    ... 
    // 先处理队列中的旧数据包
    process_send_queue();

    // 处理新数据包
    for (j=0; j<4; j++) {
        struct sockaddr_in addr = {.sin_family = AF_INET};
        addr.sin_addr.s_addr = device[j]->net_address;
        addr.sin_port = device[j]->net_port;
        
        fd_set write_fds;
        struct timeval tv;
        FD_ZERO(&write_fds);
        FD_SET(e_socket, &write_fds);
        tv.tv_sec = 0;
        tv.tv_usec = 1000;

        int select_ret = select(e_socket + 1, NULL, &write_fds, NULL, &tv);
        if (select_ret > 0 && FD_ISSET(e_socket, &write_fds)) {
            int code = sendto(e_socket, device[j]->buf, device[j]->datasize, 0, (struct sockaddr *) &addr, sizeof addr);
            if (code == -1 && errno == EAGAIN) {
                add_to_send_queue(&addr, device[j]->buf, device[j]->datasize);
            } else if (code == -1) {
                log_perror(MYNAME ": sendto error for device %d", j);
            }
        } else {
            add_to_send_queue(&addr, device[j]->buf, device[j]->datasize);
        }
    }
}

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.28 14:18:14