iotspool 原理:append-only 日志、CRC 截断恢复与退避调度
声明: iotspool 是一个开源项目,作者为 Vanderhell。本文是阅读该项目源码和文档后整理的学习笔记,用于理解嵌入式持久化消息队列的工程实现方式。本文作者不是该项目的开发者,未参与该项目的任何代码贡献。 文中所有工程细节均来自对开源代码的分析,不代表本文作者的设计决策。
项目仓库:github.com/Vanderhell/iotspool
上一篇讲了 iotspool 怎么用。本文拆解内部三件事:记录怎么编码、掉电怎么恢复、Broker 不可达时怎么退避。源码体量不大,核心逻辑都在 src/spool.c 和 src/record.c 两个文件里。
记录格式:ENQ 与 ACK 两种定长头
磁盘日志只有两种记录类型,所有整数小端序:
- magic = 0xE7:起始标志,recover 扫描时用于定位记录边界。
- type:
0x01ENQ(入队),0x02ACK(确认)。 - version = 0x01:版本字段,方便以后扩展格式,未知版本按长度跳过。
- flags:bit0 表示后面是否跟 32 字节 SHA-256,bit1 retain,bit2 qos1。
- msg_id:4 字节单调递增 ID,ACK 用它引用对应的 ENQ。
- ts_ms:入队时的单调时钟低 32 位,用于统计与退避参考。
- topic_len / payload_len:后跟 topic 和 payload 原始字节,无 null 终止符。
- crc32:最后 4 字节,覆盖整条记录除自身外所有字节。
ENQ 是变长记录:头 20 字节 + topic + payload + 可选 SHA-256 + 4 字节 CRC。ACK 固定 12 字节(4 字节头 + 4 字节 msg_id + 4 字节 CRC)。
追加写与 CRC 截断恢复
核心设计选择是只追加、不原地修改。
入队时,把 ENQ 记录拼到 scratch 缓冲区、算 CRC,调用 store.append() 写到文件尾,再调 store.sync()(POSIX 上就是 fsync)。sync 返回成功才算入队完成;如果在 append 过程中断电,这条记录只写了一半,下次启动会被识别为脏尾。
发送成功后不删除 ENQ,直接在文件尾再追加一条 ACK 记录指向同一个 msg_id。整个写入路径始终是 append-only,实现简单且掉电安全。
恢复过程(iotspool_recover)是顺序扫一遍整个 store 文件:
- 读 superblock,校验 magic/version/generation。
- 循环读下一条记录:根据 type 判断是 ENQ 还是 ACK,读固定头,算出整个记录长度(含 payload 和 CRC)。
- 若剩余字节不够读完整记录,说明写到一半断电了,
truncate_to当前位置,停止。 - 读完整记录后校验 CRC32,CRC 不对同样截断到当前位置,停止。
- CRC 正确的 ENQ 放进 RAM 索引,标记为 pending;ACK 把对应 msg_id 从索引里移除。
- 扫到文件正常结束,恢复完成。
截断到"最后一条完整记录的尾部"是关键:这保证了日志的一致性视图永远只包含已经完整 sync 过的记录,半截的永远不会被 replay。ACK 丢了最坏情况是消息重发一次(at-least-once),不会重复投递同一条多次,因为 ACK 也是 append 的。
RAM 环形索引
所有 ENQ 记录的磁盘位置不会常驻内存。RAM 里只保存一个 iotspool_entry_t 环形数组:每条 entry 包含 msg_id、generation、record_offset、record_len、topic_len、payload_len、qos、retain。
record_offset/record_len:让取出消息时直接read_at到 scratch 缓冲区,零拷贝解析。head/tail:环形队列指针,新 ENQ 追加在 tail,确认后从 head 附近移除。- 不存 topic/payload 本体,RAM 占用只跟
max_pending_msgs相关,跟消息大小无关。
这就是 STM32 上开 64 条 pending 也只占几百字节 RAM 的原因:真正的数据只在要发送的那一刻才读进 scratch。
compact:重写消除 ACK 碎片
一直 append 会让 store 文件越滚越大:已经被 ACK 的 ENQ 和对应的 ACK 记录都成了历史垃圾。当 store 大小达到 max_store_bytes 的某个阈值时,库会触发 iotspool_compact():
- 在 scratch 或工作缓冲区里重建一份新的 superblock。
- 遍历当前 RAM 索引里所有 pending 的 ENQ,用
read_at从原文件读出完整记录,直接追加到新镜像。 - 新镜像里只包含未被 ACK 的 ENQ,没有 ACK 记录,没有已确认的历史消息。
- 调用
store.replace()原子替换整个文件(POSIX 上通常是写临时文件 +rename)。 - generation 加 1,索引里的 entry 重新指向新文件中的 offset。
compact 完成后,store 大小回到"所有 pending 消息的体积之和",不会无限增长。drop_oldest_on_full=true 时,入队发现空间不够会先尝试丢最老的 pending 再 compact。
退避算法:Full Jitter
发送失败时(iotspool_on_publish_fail 或新 API 的 iotspool_publish_failed),库使用 AWS 架构博客推广的 Full Jitter 指数退避:
sleep = random_between(0, cap)
cap = min(cap * 2, max_retry_ms)初始 cap 是 min_retry_ms,每次失败翻倍直到 max_retry_ms 封顶。每次成功 ACK 后 cap 重置回 min_retry_ms。
加随机抖动是为了避免经典的"雪崩":一组设备同时断线、同时等到固定退避窗口到期、同时重连把 Broker 打死。Full Jitter 让重试时间均匀分布在 [0, cap],打散峰值。
iotspool_peek_ready(now_ms) 会比较当前时间和 retry_deadline_ms,没到点就返回 IOTSPOOL_ENOTFOUND,调用方的主循环可以去做别的事,不用自己算等待时间。
存储后端 vtable
库本身不做任何系统调用,所有 I/O 走 iotspool_store_t 的 6 个回调:
| 回调 | 职责 |
|---|---|
append(ctx, data, len) |
把数据写到当前文件尾 |
read_at(ctx, off, out, cap, *out_len) |
从指定偏移读数据 |
sync(ctx) |
落盘(fsync 或 Flash 擦写确认) |
size_bytes(ctx) |
返回当前 store 大小 |
truncate_to(ctx, new_size) |
恢复时截掉脏尾 |
replace(ctx, data, len) |
compact 时原子替换整个文件 |
POSIX 后端 store_posix.c 把这 6 个回调映射到 write/pread/fsync/fstat/ftruncate/tmpfile+rename。ESP-IDF 直接复用 POSIX 后端走 VFS。裸机上只要把这几个对接到你自己的 Flash 驱动、LittleFS、或 FAL 分区即可,核心代码不用改一行。
并发模型
库默认假设调用方自己序列化所有公开 API 调用(通常是单任务主循环 + 事件队列的嵌入式模型)。如果网络任务和采样任务会并发访问 spool,在 cfg 里填 lock/unlock 两个回调(FreeRTOS 里就是 xSemaphoreTakeRecursive / xSemaphoreGiveRecursive),核心在改内部状态和调 store 回调前后自动加锁。
锁回调必须不可重入。store 回调内部不能再调任何 iotspool API,否则死锁。ISR 里调用必须让锁回调返回 IOTSPOOL_EBUSY,或者只在 ISR 里发事件、真正 enqueue 放到任务里做。
源码阅读入口
文件极少,读的顺序推荐:
include/iotspool.h:公开 API、错误码、配置结构、entry/inflight 结构。src/record.c:记录编解码、CRC32、可选 SHA-256。src/spool.c:init/recover/enqueue/peek/ack/compact 主流程。src/backoff.c:Full Jitter 实现,几十行。src/store_posix.c:POSIX 后端,作为移植参考。tests/test_main.c:覆盖掉电截断、recover、backoff 等关键路径的单测。
设计上的三件事值得记下来:
- 写入路径永远是 append,不原地修改,CRC 覆盖整条记录,掉电只需要截尾部。
- RAM 里只存位置元数据,真正 payload 按需读,内存占用与消息大小解耦。
- compact 走整文件替换,不做段回收,逻辑短、实现好验证,跟 ENQ/ACK 双记录模型天然契合。