Iinux的双CQE通知以安全回收Buffer原理剖析
前言
本文旨在记录近期研读Java源码的学习心得与疑难问题。由于个人理解水平有限,文中内容难免存在疏漏,恳请读者不吝指正。
双CQE通知以安全回收Buffer原理剖析
1. 架构原理:双 CQE 生命周期与无锁状态机
在 Linux 6.0+ 的 IORING_OP_SEND_ZC 架构中,内核将传统套接字零拷贝复杂同步逻辑重构为基于 io_uring CQE (Completion Queue Event) 的异步状态机。
┌─────────────────────────────────────────────────────────┐
│ 应用程序提交 SQE │
│ io_uring_prep_send_zc(sqe, fd, buf...) │
└────────────────────────────┬────────────────────────────┘
│
▼
【状态: BUF_STATE_SUBMITTED】
(物理内存被 GUP 钉扎,严禁改写)
│
┌───────────────────────┴───────────────────────┐
│ │
▼ (内核成功入队) ▼ (内核直接拒绝/出错)
【产生第一 CQE (Syscall)】 【仅产生单一终结 CQE】
(cqe->flags & IORING_CQE_F_MORE = True) (cqe->flags & MORE = False)
(cqe->res = 发送字节数或错误码) (cqe->res = 错误码)
│ │
▼ │
【状态: BUF_STATE_WAIT_NOTIF】 │
(应用层获知发送字节数,但 DMA 未完成) │
│ │
▼ (网卡 DMA 完成, SKB 释放) │
【产生第二 CQE (Notification)】 │
(cqe->flags & IORING_CQE_F_MORE = False) │
(cqe->res = 0) │
│ │
└───────────────────────┬───────────────────────┘
│
▼
【状态: BUF_STATE_FREE】
(无锁回收,可安全写改或 free)
双 CQE 的精确判定契约
- 第一 CQE(传输提交结果)
cqe->res:表示推入 TCP 发送队列的字节数(若小于 0 则表示系统调用错误,如-EAGAIN或-ECONNRESET)。cqe->flags:当且仅当内核成功接管零拷贝内存并预计后续会有 DMA 释放通知时,内核会置位IORING_CQE_F_MORE。
- 第二 CQE(DMA 完成与内存释放通知)
cqe->res:固定为0(或特定的通知标记)。cqe->flags:**清零IORING_CQE_F_MORE**。收到该 CQE 代表底层 SKB 引用计数降低为 0,物理页已完全脱离网络栈。
2. 基于 liburing 的完整 C 语言示例代码
以下代码演示了基于 liburing 的 IORING_OP_SEND_ZC 零拷贝发送,包含非阻塞 Socket 对建立、Buffer 内存池设计、无锁双 CQE 解析及异常情况下的安全回收逻辑。
#define _GNU_SOURCE
#include <stdio.h>
#include <stdlib.h>
#include <string.h>
#include <unistd.h>
#include <errno.h>
#include <stdbool.h>
#include <fcntl.h>
#include <sys/socket.h>
#include <netinet/in.h>
#include <netinet/tcp.h>
#include <liburing.h>
#define RING_ENTRIES 64
#define BUF_SIZE (64 * 1024) // 64KB 大包,确保正向零拷贝收益
#define POOL_CAPACITY 8
// ============================================================================
// 1. 数据结构与状态机定义
// ============================================================================
typedef enum {
BUF_STATE_FREE = 0, // 空闲:应用层拥有完整控制权,可写改或释放
BUF_STATE_SUBMITTED, // 已提交:内核托管中,严禁修改/释放
BUF_STATE_WAIT_NOTIF // 已收到第一 CQE:发送字节数已确认,等待 DMA 完成通知
} buf_state_t;
// 缓冲区上下文节点 (携带自描述元数据)
typedef struct {
uint32_t buf_id; // 缓冲区唯一标识
buf_state_t state; // 当前状态机位置
size_t len; // 拟发送数据长度
ssize_t bytes_sent; // 实际完成发送的字节数
char *data; // 实际数据内存指针 (按页对齐)
} tx_buffer_t;
// 内存池对象 (无锁单线程环形回收)
typedef struct {
tx_buffer_t buffers[POOL_CAPACITY];
tx_buffer_t* free_stack[POOL_CAPACITY];
int top;
} buffer_pool_t;
// ============================================================================
// 2. 内存池初始化与管理逻辑
// ============================================================================
static buffer_pool_t* init_buffer_pool(void) {
buffer_pool_t *pool = calloc(1, sizeof(buffer_pool_t));
if (!pool) return NULL;
pool->top = -1;
for (int i = 0; i < POOL_CAPACITY; i++) {
pool->buffers[i].buf_id = i + 1000;
pool->buffers[i].state = BUF_STATE_FREE;
pool->buffers[i].len = BUF_SIZE;
// 按照 4096 字节对齐分配物理页,提升 GUP (Get User Pages) 效率
if (posix_memalign((void**)&pool->buffers[i].data, 4096, BUF_SIZE) != 0) {
perror("posix_memalign failed");
exit(EXIT_FAILURE);
}
// 压入空闲栈
pool->free_stack[++pool->top] = &pool->buffers[i];
}
return pool;
}
static tx_buffer_t* alloc_buffer(buffer_pool_t *pool) {
if (pool->top < 0) return NULL; // 池已耗尽
tx_buffer_t *buf = pool->free_stack[pool->top--];
buf->state = BUF_STATE_FREE;
buf->bytes_sent = 0;
return buf;
}
static void free_buffer(buffer_pool_t *pool, tx_buffer_t *buf) {
// 使用 C11/GCC 内存屏障保证状态写入前 Buffer 数据改写已完成
__atomic_store_n(&buf->state, BUF_STATE_FREE, __ATOMIC_RELEASE);
pool->free_stack[++pool->top] = buf;
}
// 设置非阻塞 Socket
static int set_nonblocking(int fd) {
int flags = fcntl(fd, F_GETFL, 0);
if (flags < 0) return -1;
return fcntl(fd, F_SETFL, flags | O_NONBLOCK);
}
// ============================================================================
// 3. 核心:双 CQE 无锁解析与回收逻辑
// ============================================================================
static void process_cqe_events(struct io_uring *ring, buffer_pool_t *pool) {
struct io_uring_cqe *cqe;
unsigned head;
unsigned count = 0;
// 【无锁批量轮询】:直接读取共享内存 CQ 环,零系统调用开销
io_uring_for_each_cqe(ring, head, cqe) {
count++;
// 从 user_data 还原缓冲区描述符指针
tx_buffer_t *buf = (tx_buffer_t *)(uintptr_t)io_uring_cqe_get_data64(cqe);
if (!buf) {
continue; // 非 zero-copy 操作或无句柄事件
}
bool has_more = (cqe->flags & IORING_CQE_F_MORE) != 0;
if (has_more) {
// ----------------------------------------------------------------
// 情况 A:收到第一 CQE (提交结果通知)
// ----------------------------------------------------------------
if (cqe->res >= 0) {
buf->bytes_sent = cqe->res;
buf->state = BUF_STATE_WAIT_NOTIF;
printf("[CQE-1 成功] Buffer ID: %u | 传输字节: %d | 标志: IORING_CQE_F_MORE | 状态 -> BUF_STATE_WAIT_NOTIF\n",
buf->buf_id, cqe->res);
printf(" └─► ⚠️ 安全屏障:DMA 传输中,物理页 %p 严禁写入/回收!\n", (void*)buf->data);
} else {
// 第一 CQE 即返回错误 (如 -EAGAIN / -EPIPE)
fprintf(stderr, "[CQE-1 错误] Buffer ID: %u | 错误码: %d (%s)\n",
buf->buf_id, cqe->res, strerror(-cqe->res));
buf->bytes_sent = cqe->res;
// 注意:即使报错,若带有 F_MORE,仍必须等待第二 CQE 才能回收!
buf->state = BUF_STATE_WAIT_NOTIF;
}
} else {
// ----------------------------------------------------------------
// 情况 B:收到终结 CQE (Notification 或者无 MORE 的单 CQE)
// ----------------------------------------------------------------
if (buf->state == BUF_STATE_WAIT_NOTIF) {
printf("[CQE-2 完成] Buffer ID: %u | DMA 释放完成 | 状态 -> BUF_STATE_FREE\n", buf->buf_id);
printf(" └─► ✅ 内存屏障解除:物理页 %p 已从内核网络栈脱离,安全的归还内存池!\n", (void*)buf->data);
} else if (buf->state == BUF_STATE_SUBMITTED) {
// 内核由于严重错误未置位 F_MORE 直接终结请求 (未产生第一 CQE)
printf("[CQE-单终结] Buffer ID: %u | 直接终结 (res=%d) | 状态 -> BUF_STATE_FREE\n",
buf->buf_id, cqe->res);
}
// 【安全回收点】:无锁归还内存池
free_buffer(pool, buf);
}
}
// 批量更新 CQ 环头指针,通知内核事件已消费
if (count > 0) {
io_uring_cq_advance(ring, count);
}
}
// ============================================================================
// 4. 主程序流程
// ============================================================================
int main(void) {
int fds[2];
struct io_uring ring;
buffer_pool_t *pool;
// 1. 创建 Unix 域套接字对用于测试
if (socketpair(AF_UNIX, SOCK_STREAM, 0, fds) < 0) {
perror("socketpair failed");
exit(EXIT_FAILURE);
}
set_nonblocking(fds[0]);
set_nonblocking(fds[1]);
// 放大套接字发送/接收缓冲区,避免因 Socket 满了退化
int sndbuf_size = 1024 * 1024;
setsockopt(fds[0], SOL_SOCKET, SO_SNDBUF, &sndbuf_size, sizeof(sndbuf_size));
setsockopt(fds[1], SOL_SOCKET, SO_RCVBUF, &sndbuf_size, sizeof(sndbuf_size));
// 2. 初始化 io_uring
if (io_uring_queue_init(RING_ENTRIES, &ring, 0) < 0) {
perror("io_uring_queue_init failed");
exit(EXIT_FAILURE);
}
// 3. 初始化内存池
pool = init_buffer_pool();
// 4. 从池中分配 Buffer 并填入数据
tx_buffer_t *buf = alloc_buffer(pool);
if (!buf) {
fprintf(stderr, "Buffer 分配失败\n");
exit(EXIT_FAILURE);
}
memset(buf->data, 'Z', BUF_SIZE);
// 5. 准备并提交零拷贝发送请求 (IORING_OP_SEND_ZC)
struct io_uring_sqe *sqe = io_uring_get_sqe(&ring);
if (!sqe) {
fprintf(stderr, "SQE 获取失败\n");
exit(EXIT_FAILURE);
}
// 使用 liburing 准备 Send Zero-Copy 算子
io_uring_prep_send_zc(sqe, fds[0], buf->data, buf->len, 0, 0);
// 将 Buffer 句柄指针直接绑定到 user_data
io_uring_sqe_set_data64(sqe, (uintptr_t)buf);
// 标记状态为已提交
buf->state = BUF_STATE_SUBMITTED;
printf("==> 提交 IORING_OP_SEND_ZC 请求 (Buffer ID: %u, Size: %zu 字节)...\n", buf->buf_id, buf->len);
io_uring_submit(&ring);
// 6. 循环等待并无锁解析双 CQE
int notif_count = 0;
while (notif_count < 2) {
// 等待至少 1 个 CQE 到达
io_uring_wait_cqe(&ring, NULL);
// 消费解析当前 CQ 环中的所有事件
process_cqe_events(&ring, pool);
// 如果 Buffer 已经被安全回收归还,说明双 CQE 流程已走完
if (pool->top >= 0 && pool->free_stack[pool->top] == buf) {
notif_count = 2; // 退出循环
}
}
// 清理资源
io_uring_queue_exit(&ring);
close(fds[0]);
close(fds[1]);
printf("==> 测试完成,Buffer 安全回收成功。\n");
return 0;
}
3. 代码关键节点与源码级机制深度拆解
3.1 内存元数据与 user_data 的无锁绑定机制
在提交 SQE 时,程序通过以下调用:
io_uring_sqe_set_data64(sqe, (uintptr_t)buf);
将应用层的 tx_buffer_t 结构体指针赋值给 sqe->user_data。
- 内核透传原理:内核在处理
IORING_OP_SEND_ZC时,分配的通知上下文struct io_notif_data会继承该user_data值。 - 双 CQE 共享关联:无论是第一 CQE(发送结果)还是第二 CQE(DMA 释放通知),内核在向共享 CQ 环压入
struct io_uring_cqe时,**均会原封不动地复制相同的user_data**。这使得应用层在处理任何一个 CQE 时,均可通过 O ( 1 ) O(1) O(1) 的指针转换恢复上下文:
tx_buffer_t *buf = (tx_buffer_t *)(uintptr_t)io_uring_cqe_get_data64(cqe);
3.2 IORING_CQE_F_MORE 标志位的无锁状态切换
在 process_cqe_events 函数中,核心判定在于逻辑分流:
bool has_more = (cqe->flags & IORING_CQE_F_MORE) != 0;
状态机转化表
| 触发事件 | cqe->res 值 | IORING_CQE_F_MORE | Buffer 转换后状态 | 内存回收动作 |
|---|---|---|---|---|
| CQE 1: 传输提交成功 | > 0 > 0 >0 (发送字节数) | Set (1) | BUF_STATE_WAIT_NOTIF | 禁止回收 (网卡 DMA 尚在读取) |
| CQE 1: 传输提交失败 | < 0 < 0 <0 (错误码) | Set (1) | BUF_STATE_WAIT_NOTIF | 禁止回收 (仍需等第二 CQE 解锁) |
| CQE 2: DMA 释放通知 | = 0 = 0 =0 | Cleared (0) | BUF_STATE_FREE | 安全回收 (归还无锁内存池) |
| 异常: 单一终结 CQE | < 0 < 0 <0 (拒绝请求) | Cleared (0) | BUF_STATE_FREE | 安全回收 (直接归还内存池) |
3.3 无锁(Lock-Free)内存安全屏障设计
在传统多线程网络库中,内存回收往往需要争抢 pthread_mutex。而在 io_uring 模式下,结合单线程 Event-Loop,可通过轻量级 C11/GCC 内存屏障保障绝对的写屏障顺序:
static void free_buffer(buffer_pool_t *pool, tx_buffer_t *buf) {
// __ATOMIC_RELEASE 保证在 state 更改为 FREE 之前,所有对于 buf->data 的读写已全部完成
__atomic_store_n(&buf->state, BUF_STATE_FREE, __ATOMIC_RELEASE);
pool->free_stack[++pool->top] = buf;
}
- 写改防护:只要
buf->state处于BUF_STATE_SUBMITTED或BUF_STATE_WAIT_NOTIF,应用层任何工作线程试图获取或修改buf->data都将被拒绝。 - CPU 乱序重排防护:
__ATOMIC_RELEASE阻止编译器和 CPU 将后续归还内存池的操作重排到 DMA 完成之前,从而在硬件层面杜绝了 Use-After-Free 或 Data Corruption 风险。
4. 生产级场景下的异常处理与边缘边界
在实际高性能网络编程中,应用层还需要处理以下三种复杂边缘场景:
4.1 短发送(Short Sends)与多次提交
由于 Socket 发送缓冲区满了,IORING_OP_SEND_ZC 的第一 CQE 返回的 cqe->res 可能是部分字节数(例如请求
64
KB
64\,\text{KB}
64KB,实际仅发送
16
KB
16\,\text{KB}
16KB)。
- 处理规则:
- 第一 CQE 返回
res = 16384且带IORING_CQE_F_MORE。 - 必须记录
buf->bytes_sent = 16384。 - 注意:该 16 KB 16\,\text{KB} 16KB 对应的物理页仍被内核 Pin 住,直到第二 CQE 达到前,整块 64 KB 64\,\text{KB} 64KB 内存依然不可写改。
- 剩余 48 KB 48\,\text{KB} 48KB 如果需要再次发送,必须使用新的偏移量分配新 SQE 提交,不能直接覆盖原 Buffer!
4.2 网络连接异常断开(EPIPE / ECONNRESET)
当客户端强制断开连接时,提交零拷贝发送可能会直接报错。
- 处理规则:
内核网络栈可能依然会触发双 CQE 流程: - CQE 1:
cqe->res = -EPIPE,flags包含IORING_CQE_F_MORE。 - CQE 2:
cqe->res = 0,flags不带IORING_CQE_F_MORE。
解析逻辑必须确保:哪怕 CQE 1 已经返回了负数错误码,只要IORING_CQE_F_MORE被置位,就绝不能提前释放内存,必须死守第二 CQE 的到来,否则将引发内核网卡 DMA 访问已释放内存的 Kernel Panic!
4.3 零拷贝退化为 Copy 场景
当发送数据包小于内核阈值(或遭遇 Copy-On-Write 页面)时,内核内部可能会将零拷贝退化为常规内核拷贝。
此时,内核仍然会遵循契约投递双 CQE(或直接在第一 CQE 后快速投递第二 CQE)。应用层无需关注内核内部是否真的使用了零拷贝,只需统一按照 IORING_CQE_F_MORE 标志位判定即可,使得应用层与内核底层实现彻底解耦。
转载自 CSDN-专业IT技术社区
原文链接:https://blog.csdn.net/wuminyu/article/details/166634161



