使用 epoll + 非阻塞I/O + 边缘触发(ET)模式 的TCP回声服务器示例

核心思路

  1. 创建监听socket,设置为非阻塞模式,并绑定到端口。

  2. 创建epoll实例(epfd)。

  3. 将监听socket添加到epoll实例中,监听可读事件(EPOLLIN),并设置为边缘触发模式(EPOLLET)。

  4. 进入主循环,调用epoll_wait等待事件。

  5. 处理事件:

    • 监听socket可读:说明有新连接到来。循环调用accept直到没有新连接为止(ET模式要求)。

    • 客户端socket可读:说明有数据到来。循环读取数据(直到读完,遇到EAGAIN),然后将数据回写给客户端。

    • 客户端socket可写:本例中,我们通常在需要回写数据时才监听可写事件,这是一个常见的优化策略。但为了简单起见,本例在连接建立后始终监听可写事件。在实际项目中,你可能需要更精细地管理可写事件的监听。

  6. 处理错误和关闭连接。

代码示例 (C语言)

#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include <arpa/inet.h>
#include <sys/socket.h>
#include <sys/epoll.h>
#include <fcntl.h>
#include <errno.h>

#define MAX_EVENTS 1024
#define BUFFER_SIZE 1024
#define PORT 8080

// 设置文件描述符为非阻塞模式
int set_nonblocking(int fd) {
    int flags = fcntl(fd, F_GETFL, 0);
    if (flags == -1) {
        perror("fcntl F_GETFL");
        return -1;
    }
    if (fcntl(fd, F_SETFL, flags | O_NONBLOCK) == -1) {
        perror("fcntl F_SETFL");
        return -1;
    }
    return 0;
}

// 添加一个fd到epoll实例,并指定监听的事件类型
void add_epoll_fd(int epoll_fd, int fd, uint32_t events) {
    struct epoll_event ev;
    ev.events = events;
    ev.data.fd = fd;
    if (epoll_ctl(epoll_fd, EPOLL_CTL_ADD, fd, &ev) == -1) {
        perror("epoll_ctl: add");
        exit(EXIT_FAILURE);
    }
}

// 修改epoll实例中已存在的fd所监听的事件
void mod_epoll_fd(int epoll_fd, int fd, uint32_t events) {
    struct epoll_event ev;
    ev.events = events;
    ev.data.fd = fd;
    if (epoll_ctl(epoll_fd, EPOLL_CTL_MOD, fd, &ev) == -1) {
        perror("epoll_ctl: mod");
        exit(EXIT_FAILURE);
    }
}

// 从epoll实例中删除一个fd
void del_epoll_fd(int epoll_fd, int fd) {
    if (epoll_ctl(epoll_fd, EPOLL_CTL_DEL, fd, NULL) == -1) {
        perror("epoll_ctl: del");
        exit(EXIT_FAILURE);
    }
}

int main() {
    int listen_sock, conn_sock, epoll_fd, nfds;
    struct sockaddr_in server_addr, client_addr;
    socklen_t client_len = sizeof(client_addr);
    struct epoll_event events[MAX_EVENTS];
    char buffer[BUFFER_SIZE];

    // 1. 创建监听socket
    listen_sock = socket(AF_INET, SOCK_STREAM, 0);
    if (listen_sock == -1) {
        perror("socket");
        exit(EXIT_FAILURE);
    }

    // 设置SO_REUSEADDR选项,避免重启后地址被占用错误
    int optval = 1;
    setsockopt(listen_sock, SOL_SOCKET, SO_REUSEADDR, &optval, sizeof(optval));

    // 2. 绑定地址和端口
    memset(&server_addr, 0, sizeof(server_addr));
    server_addr.sin_family = AF_INET;
    server_addr.sin_addr.s_addr = INADDR_ANY;
    server_addr.sin_port = htons(PORT);

    if (bind(listen_sock, (struct sockaddr*)&server_addr, sizeof(server_addr)) == -1) {
        perror("bind");
        close(listen_sock);
        exit(EXIT_FAILURE);
    }

    // 3. 开始监听
    if (listen(listen_sock, SOMAXCONN) == -1) {
        perror("listen");
        close(listen_sock);
        exit(EXIT_FAILURE);
    }

    printf("Server listening on port %d...\n", PORT);

    // 4. 创建epoll实例
    epoll_fd = epoll_create1(0);
    if (epoll_fd == -1) {
        perror("epoll_create1");
        close(listen_sock);
        exit(EXIT_FAILURE);
    }

    // 5. 将监听socket设置为非阻塞并添加到epoll,监听可读事件,使用边缘触发模式
    set_nonblocking(listen_sock);
    // EPOLLIN | EPOLLET 表示监听可读事件并使用边缘触发
    add_epoll_fd(epoll_fd, listen_sock, EPOLLIN | EPOLLET); 

    // 主循环
    while (1) {
        // 6. 等待事件发生
        nfds = epoll_wait(epoll_fd, events, MAX_EVENTS, -1); // -1 表示无限期阻塞
        if (nfds == -1) {
            perror("epoll_wait");
            break;
        }

        for (int i = 0; i < nfds; i++) {
            int current_fd = events[i].data.fd;

            // 7. 处理新连接 (监听socket可读)
            if (current_fd == listen_sock) {
                printf("Listening socket is readable, accepting new connections...\n");
                // ET模式要求必须循环accept直到没有新连接为止
                while (1) {
                    conn_sock = accept(listen_sock, (struct sockaddr*)&client_addr, &client_len);
                    if (conn_sock == -1) {
                        // 如果没有更多新连接了,就跳出循环
                        if (errno == EAGAIN || errno == EWOULDBLOCK) {
                            break; // 已经accept完所有新连接
                        } else {
                            perror("accept");
                            break;
                        }
                    }

                    // 打印客户端信息
                    char client_ip[INET_ADDRSTRLEN];
                    inet_ntop(AF_INET, &client_addr.sin_addr, client_ip, INET_ADDRSTRLEN);
                    printf("Accepted new connection from %s:%d, fd: %d\n",
                           client_ip, ntohs(client_addr.sin_port), conn_sock);

                    // 设置新连接的socket为非阻塞模式
                    set_nonblocking(conn_sock);

                    // 将新连接的socket添加到epoll,监听可读和可写事件,使用边缘触发模式
                    // 注意:实际项目中,通常只在需要写入数据时才监听可写事件(EPOLLOUT)以避免 busy loop
                    uint32_t conn_events = EPOLLIN | EPOLLOUT | EPOLLET | EPOLLRDHUP;
                    add_epoll_fd(epoll_fd, conn_sock, conn_events);

                    // 可以在这里发送欢迎信息等...
                    // const char *welcome_msg = "Welcome to the epoll echo server!\n";
                    // send(conn_sock, welcome_msg, strlen(welcome_msg), 0);
                }
            }
            // 8. 处理客户端socket的可读事件
            else if (events[i].events & EPOLLIN) {
                printf("Client fd %d is readable.\n", current_fd);
                // ET模式要求必须循环read直到读完
                while (1) {
                    ssize_t count = read(current_fd, buffer, BUFFER_SIZE);
                    if (count == -1) {
                        // 数据读完了,或者暂时没数据了
                        if (errno == EAGAIN || errno == EWOULDBLOCK) {
                            break; // 等待下一次EPOLLIN事件
                        } else {
                            perror("read");
                            close(current_fd);
                            del_epoll_fd(epoll_fd, current_fd);
                            break;
                        }
                    } else if (count == 0) {
                        // 客户端关闭了连接
                        printf("Client fd %d closed connection.\n", current_fd);
                        close(current_fd);
                        del_epoll_fd(epoll_fd, current_fd);
                        break;
                    } else {
                        // 成功读到数据
                        printf("Received %zd bytes from fd %d: %.*s",
                               count, current_fd, (int)count, buffer);
                        // 注意:这里只是简单地把数据放回buffer,由可写事件负责发送回去
                        // 在实际应用中,你可能需要将数据与socket fd关联起来(例如使用一个结构体)
                        // 这里为了简单,我们假设可写事件会立刻将数据发回。
                    }
                }
            }
            // 9. 处理客户端socket的可写事件 (有空间可以发送数据)
            else if (events[i].events & EPOLLOUT) {
                // 本例中,我们简单地将收到的数据回显。
                // 注意:这是一个过于简化的例子。在实际中,你需要管理每个连接的发送缓冲区。
                // 这里我们只是象征性地发送一个回应。
                const char *response = "Echo: ";
                // 假设buffer中存储着最后收到的数据(这在实际并发环境中是不对的!)
                // 正确做法是为每个连接维护一个发送缓冲区队列。
                send(current_fd, response, strlen(response), 0);
                send(current_fd, buffer, strlen(buffer), 0); // 注意:这里可能有问题!
            }
            // 10. 处理错误或挂起事件 (EPOLLERR | EPOLLHUP)
            else if ( (events[i].events & EPOLLERR) || (events[i].events & EPOLLHUP) ) {
                printf("Error or hangup on fd %d. Closing.\n", current_fd);
                close(current_fd);
                del_epoll_fd(epoll_fd, current_fd);
            }
        }
    }

    // 清理
    close(epoll_fd);
    close(listen_sock);
    return 0;
}

注意:

  1. 数据管理:代码有一个严重问题buffer 是全局共享的。在并发情况下,多个客户端会覆盖彼此的缓冲区。正确的做法是为每个连接创建一个数据结构(struct),包含其独有的读缓冲区和写缓冲区。

  2. 可写事件处理:监听 EPOLLOUT 要小心。如果发送缓冲区一直有空位,epoll_wait 会一直返回可写事件,导致 busy loop。最佳实践是:只在有数据要发送时才监听可写事件,数据发送完后立即取消监听。

  3. 错误处理:示例中的错误处理比较简单,生产环境需要更健壮。

  4. ET模式:务必记住使用非阻塞I/O并循环读/写直到 EAGAIN

Logo

有“AI”的1024 = 2048,欢迎大家加入2048 AI社区

更多推荐