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 源码里没有这些,看到别懵。

代码块C · 49 行收起展开
// 基于本地 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 最终都落到这里:

代码块C · 148 行收起展开
// 基于本地 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 时双删:

代码块C · 84 行收起展开
// 基于本地 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。
写完后 xaddCommandrewriteClientCommandArgument 把命令里的 * 替换成真实生成的 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 拿到消费者,
然后 streamReplyWithRangegroup->last_id 之后开始迭代(streamIteratorGetID 负责从 rax key 解出 master_id 再加 delta 还原每条 ID)。
每回复一条就 streamUpdateCGroupLastId 推进组位移,并 streamCreateNACK 建未确认记录、raxTryInsert 进组 PEL、同一指针再插进消费者自己的 PEL。

代码块JAVA · 2 行收起展开
为什么投递时就推进 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 不内置死信队列。

延伸阅读