如何解决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

