Stream
stream 源码分析(消息流)
Stream 是 Redis 5.0 加入的消息队列类型:一条条带全局递增 ID 的消息,支持消费者组、ack、阻塞读。
它解决的问题:List 做队列没有广播和重复消费能力,Pub/Sub 不落盘、掉线就丢消息。
核心思路是两层结构——rax 按 ID 做有序索引(见 rax),每个 rax 节点挂一个 listpack 批量压缩存多条消息(见 listpack);消费者组语义则全靠一张”已投递未确认”的 PEL 表实现至少一次投递。
注意:以下代码摘自本地 unstable 分支,比 7.x 多了 cgroups_ref、PEL 时间链表、IDMP(幂等生产者)等新字段,网上教程和 7.x 源码里没有这些,看到别懵。
代码块收起展开
// 基于本地 Redis 仓 (unstable), src/stream.h
typedef struct streamID {
uint64_t ms; // Unix 毫秒时间戳
uint64_t seq; // 同毫秒内的序号,ID 形如 1699000000000-0
} streamID;
typedef struct stream {
rax *rax; // 索引树:key = 消息 ID 编码成 16 字节大端,value = 一个 listpack(批量存多条消息)
uint64_t length; // 当前存活消息数(不含 tombstone)
streamID last_id; // 最后生成的 ID,新 ID 必须严格大于它,全局递增的唯一依据
streamID first_id; // 第一条非 tombstone 消息的 ID
streamID max_deleted_entry_id; // 删过的最大 ID,XINFO / lag 估算用
uint64_t entries_added; // 历史累计写入数(含已删),和 length 的差就是删除量
size_t alloc_size; // 该 stream 占用的总内存字节数
rax *cgroups; // 消费者组:组名 -> streamCG。不建组就是 NULL,按需创建省内存
rax *cgroups_ref; // 反向索引:消息 ID -> 引用它的组,XACKDEL/XDELEX 判断"还有没有组引用着"
streamID min_cgroup_last_id; // 所有组 last_id 的最小值,裁剪时判断"是否所有组都消费过了"
unsigned int min_cgroup_last_id_valid: 1;
// ... // IDMP 幂等生产者去重字段(unstable 在研特性),略
} stream;
/* 消费者组:实现"组内一条消息只给一个人 + 必须 ack"的语义 */
typedef struct streamCG {
streamID last_id; // 组已投递到的位置,XREADGROUP 用 ">" 就是从它之后取新消息
long long entries_read; // 组累计读取数的估算值,算 lag 用;有删除/SETID 时会失真需重估
rax *pel; // Pending Entries List:已投递未 ack 的消息,ID -> streamNACK,"至少一次"的核心
streamNACK *pel_time_head; // PEL 按 delivery_time 排序的双向链表头(最老),XAUTOCLAIM O(1) 拿最老待认领
streamNACK *pel_time_tail; // 链表尾(最新),重投递刷新时间时 O(1) 挪到尾部
rax *consumers; // 消费者名 -> streamConsumer
} streamCG;
typedef struct streamConsumer {
mstime_t seen_time; // 最后一次尝试动作的时间(读失败也算)
mstime_t active_time; // 最后一次真正读到/认领到消息的时间,XINFO 靠它识别僵尸消费者
sds name; // 消费者名,大小写敏感;不用注册,首次使用即创建
rax *pel; // 该消费者自己的未 ack 列表,value 与组 pel 指向同一个 streamNACK,共享不复制
} streamConsumer;
/* 未确认消息记录:谁领的、何时领的、投递过几次 */
struct streamNACK {
mstime_t delivery_time; // 最后一次投递时间,XCLAIM/XAUTOCLAIM 用它算空闲时长
uint64_t delivery_count; // 投递次数,反复处理失败的毒消息靠它识别
streamConsumer *consumer; // 当前归属的消费者
listNode *cgroup_ref_node; // 自己在 cgroups_ref 里的节点,删除时 O(1) 摘除
streamID id; // 冗余存一份 ID:走时间链表遍历时不用再从 rax key 反解
struct streamNACK *pel_prev; // 时间序双向链表指针(unstable 新增,为 XAUTOCLAIM 提速)
struct streamNACK *pel_next;
};写入路径的核心是 streamAppendItem,XADD 最终都落到这里:
代码块收起展开
// 基于本地 Redis 仓 (unstable), src/t_stream.c
/* 生成下一个 ID:时钟正常就用当前毫秒 + seq=0;时钟回拨则沿用 last_id 的毫秒,只递增 seq */
void streamNextID(streamID *last_id, streamID *new_id) {
uint64_t ms = commandTimeSnapshot();
if (ms > last_id->ms) {
new_id->ms = ms;
new_id->seq = 0;
} else {
*new_id = *last_id; // 时钟回拨也绝不后退:宁可时间戳"虚高",也要保住全局递增
streamIncrID(new_id);
}
}
/* XADD 底层实现:把一条 field-value 消息追加进 rax + listpack */
int streamAppendItem(stream *s, robj **argv, int64_t numfields, streamID *added_id, streamID *use_id, int seq_given) {
streamID id;
if (use_id) {
if (seq_given) {
id = *use_id; // 显式完整 ID:主从复制回放走这条路
} else {
if (s->last_id.ms == use_id->ms) { // "ms-*" 形式:同毫秒自动续 seq
if (s->last_id.seq == UINT64_MAX) {
errno = EDOM;
return C_ERR;
}
id = s->last_id;
id.seq++;
} else {
id = *use_id;
}
}
} else {
streamNextID(&s->last_id,&id); // "*":自动生成
}
if (streamCompareID(&id,&s->last_id) <= 0) { // 硬门槛:必须严格大于 last_id,否则 XADD 直接报错
errno = EDOM;
return C_ERR;
}
size_t totelelen = 0;
for (int64_t i = 0; i < numfields*2; i++) {
sds ele = argv[i]->ptr;
totelelen += sdslen(ele);
}
if (totelelen > STREAM_LISTPACK_MAX_SIZE) { // listpack 长度字段是 32 位,超大消息直接拒
errno = ERANGE;
return C_ERR;
}
raxIterator ri;
raxStart(&ri,s->rax);
raxSeek(&ri,"$",NULL,0); // "$" = 定位 rax 最后一个节点:追加永远发生在尾部 listpack
size_t lp_bytes = 0;
unsigned char *lp = NULL;
if (!raxEOF(&ri)) {
lp = ri.data;
lp_bytes = lpBytes(lp);
}
raxStop(&ri);
uint64_t rax_key[2];
streamID master_id;
/* 尾部节点满了(字节数或条数超限)就封存换新节点 */
if (lp != NULL) {
int new_node = 0;
size_t node_max_bytes = server.stream_node_max_bytes;
if (node_max_bytes == 0 || node_max_bytes > STREAM_LISTPACK_MAX_SIZE)
node_max_bytes = STREAM_LISTPACK_MAX_SIZE;
if (lp_bytes + totelelen >= node_max_bytes) {
new_node = 1;
} else if (server.stream_node_max_entries) {
unsigned char *lp_ele = lpFirst(lp);
int64_t count = lpGetInteger(lp_ele) + lpGetInteger(lpNext(lp,lp_ele)); // 存活+已删都算:限的是节点遍历成本
if (count >= server.stream_node_max_entries) new_node = 1;
}
if (new_node) {
lp = lpShrinkToFit(lp); // 封存前把预分配的多余内存缩掉
// ...
lp = NULL;
}
}
int flags = STREAM_ITEM_FLAG_NONE;
if (lp == NULL) {
master_id = id;
streamEncodeID(rax_key,&id); // ID 编成 128 位大端当 rax key:字典序 == 数值序,rax 范围查找才成立
/* 新建 listpack,头部写 master entry:count | deleted | num-fields | field1..N | 0 */
size_t prealloc = STREAM_LISTPACK_MAX_PRE_ALLOCATE; // 预分配一块,避免每次 XADD 都 realloc
if (server.stream_node_max_bytes > 0 && server.stream_node_max_bytes < prealloc) {
prealloc = server.stream_node_max_bytes;
}
lp = lpNew(prealloc);
lp = lpAppendInteger(lp,1);
lp = lpAppendInteger(lp,0);
lp = lpAppendInteger(lp,numfields);
for (int64_t i = 0; i < numfields; i++) {
sds field = argv[i*2]->ptr;
lp = lpAppend(lp,(unsigned char*)field,sdslen(field)); // master entry 只存 field 名,当后续消息的模板
}
lp = lpAppendInteger(lp,0);
s->alloc_size += lpBytes(lp);
raxInsert(s->rax,(unsigned char*)&rax_key,sizeof(rax_key),lp,NULL);
flags |= STREAM_ITEM_FLAG_SAMEFIELDS;
// ...
} else {
// ... 读出节点头部 master entry,逐个比对本条消息的 field 名
if (numfields == master_fields_count) {
// ...
if (i == master_fields_count) flags |= STREAM_ITEM_FLAG_SAMEFIELDS; // 全同才置位:本条不再存 field 名
}
}
/* 真实条目编码:flags | ms-diff | seq-diff | [num-fields | field...] | value... | lp-count */
size_t oldsize = lpBytes(lp);
lp = lpAppendInteger(lp,flags);
lp = lpAppendInteger(lp,id.ms - master_id.ms); // ID 只存与 master 的差值(delta 编码),小整数编码更省
lp = lpAppendInteger(lp,id.seq - master_id.seq);
if (!(flags & STREAM_ITEM_FLAG_SAMEFIELDS))
lp = lpAppendInteger(lp,numfields);
for (int64_t i = 0; i < numfields; i++) {
sds field = argv[i*2]->ptr, value = argv[i*2+1]->ptr;
if (!(flags & STREAM_ITEM_FLAG_SAMEFIELDS))
lp = lpAppend(lp,(unsigned char*)field,sdslen(field)); // SAMEFIELDS 时 field 名一个字节都不存
lp = lpAppend(lp,(unsigned char*)value,sdslen(value));
}
int64_t lp_count = numfields;
lp_count += 3;
if (!(flags & STREAM_ITEM_FLAG_SAMEFIELDS)) {
lp_count += numfields+1;
}
lp = lpAppendInteger(lp,lp_count); // 尾部记本条占几个 listpack 元素,XREVRANGE 反向遍历靠它往回跳
s->alloc_size -= oldsize;
s->alloc_size += lpBytes(lp);
if (ri.data != lp)
raxInsert(s->rax,(unsigned char*)&rax_key,sizeof(rax_key),lp,NULL); // append 可能 realloc 挪地址,必须回写指针
s->length++;
s->entries_added++;
s->last_id = id;
if (s->length == 1) s->first_id = id;
if (added_id) *added_id = id;
return C_OK;
}消费侧三个关键函数:投递时建 NACK、读取时登记 PEL、ack 时双删:
代码块收起展开
// 基于本地 Redis 仓 (unstable), src/t_stream.c
/* 消息投递给消费者的那一刻创建未确认记录 */
streamNACK *streamCreateNACK(stream *s, streamConsumer *consumer, streamID *id) {
size_t usable;
streamNACK *nack = zmalloc_usable(sizeof(*nack), &usable);
s->alloc_size += usable;
nack->delivery_time = commandTimeSnapshot();
nack->delivery_count = 1;
nack->consumer = consumer;
nack->cgroup_ref_node = NULL;
nack->id = *id;
nack->pel_prev = NULL;
nack->pel_next = NULL;
return nack;
}
/* XREAD/XRANGE/XREADGROUP 共用的读取核心;带 group 时边回复边登记 PEL */
size_t streamReplyWithRange(client *c, stream *s, streamID *start, streamID *end, size_t count, int rev, long long min_idle_time, streamCG *group, streamConsumer *consumer, int flags, streamPropInfo *spi, unsigned long *propCount) {
// ...
streamIteratorStart(&si,s,start,end,rev);
while (streamIteratorGetID(&si,&id,&numfields)) {
if (group && streamCompareID(&id,&group->last_id) > 0) {
// ...
streamUpdateCGroupLastId(s, group, &id); // 投递即推进 last_id,不等 ack:组内不会把同一条发给第二个人
propagate_last_id = 1; // last_id 变化要传播到 AOF/副本,否则切主后重复消费
}
// ... 把 ID 和 field-value 写进客户端回复
if (group && !noack) { // NOACK 模式不进 PEL:主动退化成"至多一次"
unsigned char buf[sizeof(streamID)];
streamEncodeID(buf,&id);
streamNACK *nack = streamCreateNACK(s, consumer, &id);
int group_inserted =
raxTryInsert(group->pel,buf,sizeof(buf),nack,NULL); // 乐观先插:ID 绝大多数是新的,省一次查找
if (group_inserted == 0) { // 已存在(XGROUP SETID 把 last_id 往回拨过):转给新消费者
streamFreeNACK(s,nack);
void *result;
int found = raxFind(group->pel,buf,sizeof(buf),&result);
serverAssert(found);
nack = result;
if (nack->consumer != consumer) {
raxRemove(nack->consumer->pel,buf,sizeof(buf),NULL);
nack->consumer = consumer;
raxInsert(consumer->pel,buf,sizeof(buf),nack,NULL);
}
nack->delivery_count = 1;
pelListUpdate(group, nack, cmd_time_snapshot);
} else {
raxInsert(consumer->pel,buf,sizeof(buf),nack,NULL); // 组 PEL 和消费者 PEL 挂同一个 nack 指针
nack->cgroup_ref_node = streamLinkCGroupToEntry(s, group, buf);
pelListInsertAtTail(group, nack); // 新投递时间必然最新,直接挂时间链表尾,O(1) 保持有序
}
consumer->active_time = cmd_time_snapshot;
// ...
}
decrRefCount(idarg);
arraylen++;
if (count && count == arraylen) break;
}
// ...
return arraylen;
}
/* XACK:从两个 PEL 里摘掉,"处理完成"就此成立 */
void xackCommand(client *c) {
// ...
for (int j = 3; j < c->argc; j++) {
unsigned char buf[sizeof(streamID)];
streamEncodeID(buf,&ids[j-3]);
void *result;
if (raxFind(group->pel,buf,sizeof(buf),&result)) {
streamNACK *nack = result;
pelListUnlink(group, nack); // 摘时间链表
raxRemove(group->pel,buf,sizeof(buf),NULL); // 摘组 PEL
raxRemove(nack->consumer->pel,buf,sizeof(buf),NULL); // 摘消费者 PEL:nack 记着归属,不用遍历所有消费者
streamDestroyNACK(kv->ptr, nack, buf);
acknowledged++;
server.dirty++;
keyModified(c,c->db,c->argv[1],kv,0);
}
}
addReplyLongLong(c,acknowledged);
// ...
}原理串讲
一次完整的生产消费长这样。生产端执行 XADD key * f v:xaddCommand 解析参数后调 streamAppendItem,先由 streamNextID 用当前毫秒生成 ID(时钟回拨就沿用旧毫秒只加 seq,streamCompareID 再兜底一道”必须大于 last_id”),
然后 raxSeek("$") 定位到 rax 尾节点的 listpack:没满就把消息以 delta 编码追加进去,满了就 lpShrinkToFit 封存旧节点、新建 listpack 并用本条 ID 当 rax key。
写完后 xaddCommand 调 rewriteClientCommandArgument 把命令里的 * 替换成真实生成的 ID 再传播——为什么?因为 AOF 重放和副本不能各自再生成一遍时间戳,否则主从的消息 ID 会分叉,所有依赖 ID 的 PEL、last_id 全乱。
最后 signalKeyAsReady 唤醒阻塞在 XREAD BLOCK 上的客户端。
为什么 rax key 是 128 位大端而不是直接存结构体?rax 只会按字节做字典序,把 ms 和 seq 各自转大端拼起来,字典序恰好等于数值序,XRANGE 的范围查询才能直接落在 rax 的有序遍历上;同时相邻消息的 ms 高位字节几乎相同,rax 的前缀压缩把这部分只存一份。
而一个 rax 节点里挂 listpack 批量存多条,是因为消息的 field 名高度重复:节点首条消息的 field 列表做成 master entry 当模板,后续消息若 field 完全一致就置 STREAM_ITEM_FLAG_SAMEFIELDS,只存 value,连 ID 都只存与 master 的差值,内存省一个量级。
消费端 XREADGROUP GROUP g c1 COUNT 10 STREAMS key > 进入 xreadCommand:streamLookupCG 找到组、streamLookupConsumer/streamCreateConsumer 拿到消费者,
然后 streamReplyWithRange 从 group->last_id 之后开始迭代(streamIteratorGetID 负责从 rax key 解出 master_id 再加 delta 还原每条 ID)。
每回复一条就 streamUpdateCGroupLastId 推进组位移,并 streamCreateNACK 建未确认记录、raxTryInsert 进组 PEL、同一指针再插进消费者自己的 PEL。
代码块收起展开
为什么投递时就推进 last_id 而非 ack 时?因为组内是竞争消费,位移不推进的话同一条消息会被下一个 `XREADGROUP` 再发一遍;
"消息可能丢"的风险不由位移承担,而由 PEL 承担——消费者处理完调 `XACK`,`xackCommand` 用 `raxFind` 找到 NACK 后从组 PEL、消费者 PEL、时间链表三处摘除;消费者中途崩了不 ack,NACK 留在 PEL 里,别的消费者用 XAUTOCLAIM 沿 pel_time_head 时间链表 O(1) 找到最老的超时消息接手(这条链表是 unstable 新加的,7.x 的 XAUTOCLAIM 得从头扫 PEL 的 rax)。
这套 PEL 机制就是 Stream 的至少一次投递:状态全在服务端,客户端崩溃不丢”谁欠着哪条没确认”的账。
设计取舍
- 至少一次而非精确一次:ack 前宕机必然重复投递,消费端必须幂等。和 Kafka 手动提交位移、RocketMQ 消费重试是同一个问题的同一种答案。
- ID 即位移,且位移存在服务端(
streamCG.last_id+ PEL),不像 Kafka 靠消费者自己管 offset;代价是每个组在 Redis 内存里都有一份状态。 XDEL只在 listpack 里把条目标成 tombstone,整个节点删空才真正释放内存;所以XLEN不降内存不一定降。XADD MAXLEN ~ 1000的近似裁剪只删整个 rax 节点,精确裁剪要改写 listpack,高吞吐下永远该用~。delivery_count只计数不兜底:毒消息处理(超次数转死信)要业务自己用XPENDING/XAUTOCLAIM实现,Redis 不内置死信队列。