欢迎光临
我们一直在努力

【学习笔记】io_uring 原理、执行流程与 TCP Echo Server 实现

io_uring 是什么:

io_uring 是 Linux 提供的一套异步 I/O 接口。

传统网络编程中,我们经常使用:

accept()
recv()
send()
read()
write()

这些函数的特点是:应用程序主动进入内核执行一次 I/O 操作。

后来为了处理大量连接,我们会使用 select、poll、epoll。

但 epoll 本质上解决的是:哪些 fd 现在已经准备好,可以进行 I/O?

之后我们其实还是需要进行一次手动的调用处理。

io_uring实际上是告诉内核:直接进行一次比如accept()的操作,我们不用手动操作。

大概流程是: 应用程序准备 I/O 请求。 请求放入 SQ。 内核取走请求并执行。 操作完成。 内核把完成结果放入 CQ。 应用程序从 CQ 中取出结果并处理。 根据结果继续提交新的操作。

io_uring的基本调用:

SQ:Submission Queue

SQ 是提交队列。

应用程序想让内核执行什么操作,就准备对应的 SQE,然后放入 SQ。

SQ 保存的是准备提交或者等待内核取走的 I/O 请求

SQE:Submission Queue Entry

SQ 是一个队列,而 SQE 是里面的一个具体请求。

CQ:Completion Queue

CQ 是完成队列。

内核完成一个 I/O 操作以后,会产生一个 CQE,并放到 CQ 中。

CQE:Completion Queue Entry

CQE 是某一个已经完成的 I/O 操作对应的结果。

io_uring 最核心的运行流程:

假设现在要执行:

recv(sockfd, buffer, 1024, 0);

传统方式是直接调用:

recv()

而 io_uring 会分成几个阶段。

第一阶段:获取 SQE

struct io_uring_sqe *sqe;

sqe = io_uring_get_sqe(&ring);

从 SQ 中获得一个可以填写的 SQE。

第二阶段:准备操作

io_uring_prep_recv(sqe, sockfd, buffer, 1024, 0);

告诉 SQE:

操作类型:recv
fd:sockfd
buffer:buffer
最大长度:1024
flags:0

现在只是:

准备好了一个任务

还没有真正完成 recv。

第三阶段:提交请求

io_uring_submit(&ring);

将准备好的 SQE 提交给内核。

一个重要优势是:SQ 中可以一次准备多个 SQE,然后批量提交。

所以可能一次准备:

recv A
recv B
recv C
accept
send D

然后一次 io_uring_submit()。

这也是 io_uring 可以减少频繁用户态/内核态交互的重要原因之一。SQ Polling 模式下还可以由专门的内核线程轮询 SQ,从而进一步减少提交路径上的系统调用。

第四阶段:内核执行 I/O

内核开始处理:

recv

此时应用程序不需要一直围绕这个 recv 阻塞执行自己的处理逻辑。

第五阶段:操作完成

假设收到:

hello

5 个字节。

内核产生一个 CQE。

例如:

cqe->res == 5

第六阶段:应用处理 CQE

应用程序从 CQ 中获取完成事件:

io_uring_wait_cqe(&ring, &cqe);

然后根据 CQE 判断:

完成的是谁?
完成的是什么操作?
结果是什么?

再决定下一步做什么。

io_uring 需要 user_data:

假设现在同时提交:

客户端A recv
客户端B recv
客户端C send
监听socket accept

过了一会 CQ 中出现一个 CQE。

程序必须知道:

这个CQE属于哪个fd?

它原来执行的是accept、recv还是send?

所以需要给每一个 SQE 保存自己的上下文。

因此需要一个user_data来记录

io_uring 和 epoll 的本质区别:

epoll偏向于就绪通知,后续还需要手动操作一下。

io_uring偏向于完成通知,因为内核已经完成了比如accept()的操作了。

io_uring性能好的原因:

1.共享 Ring Buffer

SQ 和 CQ 使用用户态和内核态之间共享的 ring buffer。

应用程序通过 SQ 描述请求,内核通过 CQ 返回结果。

减少了很多传统 I/O 接口中的额外数据交换和管理开销。

2.批量提交

可以:

准备很多SQE

然后:

io_uring_submit()

批量提交。

不是每一个 I/O 操作都一定要单独完成一次完整的提交路径。

3.批量获取完成事件

例如代码:

struct io_uring_cqe *cqes[128];

int nready = io_uring_peek_batch_cqe(&ring, cqes, 128);

可以一次拿很多已经完成的 CQE。

4.SQPOLL

io_uring 支持 SQ Polling 模式。

可以让内核线程轮询 SQ。

普通模式:

用户准备SQE
用户通知内核
内核处理

SQPOLL 模式则可以在适合的场景下降低通知内核的系统调用开销。

5.Registered Buffer

io_uring 可以提前把 buffer 注册给内核。

这样在大量频繁 I/O 中,可以减少反复建立内存映射、固定页面等相关开销。

6.Registered File

fd 同样可以注册到 io_uring 中。

内核可以长期持有这些资源的引用,降低大量 I/O 请求中反复处理 fd 的部分开销。

完整实现代码:

#include <stdio.h>
#include <liburing.h>
#include <netinet/in.h>
#include <string.h>
#include <unistd.h>

#define EVENT_ACCEPT0//accept完成事件
#define EVENT_READ1//recv完成事件
#define EVENT_WRITE2//send完成事件

//保存一个io_uring请求对应的信息
//因为CQE回来以后,我们需要知道这个完成事件原来是什么操作
struct conn_info {
int fd;//这个事件对应的fd
int event;//事件类型:ACCEPT/READ/WRITE
};

//创建TCP监听服务器
int init_server(unsigned short port) {

int sockfd = socket(AF_INET, SOCK_STREAM, 0);//创建TCPsocket

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

serveraddr.sin_family = AF_INET;//IPv4
serveraddr.sin_addr.s_addr = htonl(INADDR_ANY);//监听本机所有IP
serveraddr.sin_port = htons(port);//监听指定端口

if (-1 == bind(sockfd, (struct sockaddr*)&serveraddr, sizeof(struct sockaddr))) {
perror("bind");
return -1;
}

listen(sockfd, 10);//进入监听状态,backlog=10

return sockfd;
}

#define ENTRIES_LENGTH1024//io_uring队列容量
#define BUFFER_LENGTH1024//收发缓冲区大小

//向io_uring的SQ中添加一个recv请求
int set_event_recv(struct io_uring *ring, int sockfd,
void *buf, size_t len, int flags) {

struct io_uring_sqe *sqe = io_uring_get_sqe(ring);//从SQ中获取一个SQE

struct conn_info accept_info = {
.fd = sockfd,
.event = EVENT_READ,
};

io_uring_prep_recv(sqe, sockfd, buf, len, flags);//把这个SQE准备成recv操作

//把fd和事件类型保存到SQE的user_data中
//以后这个请求完成时,CQE会把user_data原样带回来
memcpy(&sqe->user_data, &accept_info, sizeof(struct conn_info));
}

//向io_uring的SQ中添加一个send请求
int set_event_send(struct io_uring *ring, int sockfd,
void *buf, size_t len, int flags) {

struct io_uring_sqe *sqe = io_uring_get_sqe(ring);//从SQ中获取一个SQE

struct conn_info accept_info = {
.fd = sockfd,
.event = EVENT_WRITE,
};

io_uring_prep_send(sqe, sockfd, buf, len, flags);//把SQE准备成send操作

//保存这个请求对应的fd和事件类型
memcpy(&sqe->user_data, &accept_info, sizeof(struct conn_info));
}

//向io_uring的SQ中添加一个accept请求
int set_event_accept(struct io_uring *ring, int sockfd, struct sockaddr *addr,
socklen_t *addrlen, int flags) {

struct io_uring_sqe *sqe = io_uring_get_sqe(ring);//从SQ中获取一个SQE

struct conn_info accept_info = {
.fd = sockfd,
.event = EVENT_ACCEPT,
};

io_uring_prep_accept(sqe, sockfd, (struct sockaddr*)addr, addrlen, flags);//准备accept操作

//保存监听fd以及事件类型
memcpy(&sqe->user_data, &accept_info, sizeof(struct conn_info));
}

int main(int argc, char *argv[]) {

unsigned short port = 9999;
int sockfd = init_server(port);//创建监听socket

struct io_uring_params params;
memset(&params, 0, sizeof(params));

struct io_uring ring;

//初始化io_uring
//ENTRIES_LENGTH表示SubmissionQueue最多允许准备多少个请求
io_uring_queue_init_params(ENTRIES_LENGTH, &ring, &params);

#if 0
struct sockaddr_in clientaddr;
socklen_t len = sizeof(clientaddr);

//传统阻塞accept
accept(sockfd, (struct sockaddr*)&clientaddr, &len);
#else

struct sockaddr_in clientaddr;
socklen_t len = sizeof(clientaddr);

//不直接调用accept
//而是向io_uring提交一个accept请求
set_event_accept(&ring, sockfd, (struct sockaddr*)&clientaddr, &len, 0);

#endif

char buffer[BUFFER_LENGTH] = {0};

while (1) {

//把SQ中准备好的请求真正提交给内核
io_uring_submit(&ring);

struct io_uring_cqe *cqe;

//等待至少有一个请求完成
io_uring_wait_cqe(&ring, &cqe);

struct io_uring_cqe *cqes[128];

//一次从CQ中获取最多128个已经完成的事件
//作用上可以暂时类比epoll_wait返回一批事件
int nready = io_uring_peek_batch_cqe(&ring, cqes, 128);

int i = 0;

for (i = 0;i < nready;i ++) {

struct io_uring_cqe *entries = cqes[i];

struct conn_info result;

//把之前放在SQE->user_data中的信息取回来
memcpy(&result, &entries->user_data, sizeof(struct conn_info));

if (result.event == EVENT_ACCEPT) {

//上一个accept完成之后
//马上再准备一个新的accept
//这样服务器可以继续接收其他客户端连接
set_event_accept(&ring, sockfd, (struct sockaddr*)&clientaddr, &len, 0);

//accept的返回值放在CQE->res中
//成功时就是新连接的connfd
int connfd = entries->res;

//给刚连接进来的客户端准备一个recv操作
set_event_recv(&ring, connfd, buffer, BUFFER_LENGTH, 0);

} else if (result.event == EVENT_READ) {

//recv操作的返回值
//>0表示读取到的字节数
//=0表示客户端关闭连接
//<0表示发生错误
int ret = entries->res;

if (ret == 0) {

//客户端断开连接
close(result.fd);

} else if (ret > 0) {

//收到多少字节
//就发送多少字节回客户端
set_event_send(&ring, result.fd, buffer, ret, 0);
}

} else if (result.event == EVENT_WRITE) {

//send操作完成
int ret = entries->res;

//发送完成以后
//继续给这个客户端准备下一次recv
set_event_recv(&ring, result.fd, buffer, BUFFER_LENGTH, 0);
}
}

//告诉io_uring
//前面nready个CQE我们已经处理完了
//CQ可以释放这些位置继续使用
io_uring_cq_advance(&ring, nready);
}
}

零声社区资源链接:0voice · GitHub

赞(0)
未经允许不得转载:171主机测评 » 【学习笔记】io_uring 原理、执行流程与 TCP Echo Server 实现
分享到: 更多 (0)

评论 抢沙发

  • 昵称 (必填)
  • 邮箱 (必填)
  • 网址