Redis中Stream的基础数据结构
消息ID
typedef struct streamID {uint64_t ms;uint64_t seq;} streamID;
Redis使用streamID
这个数据结构来表示消息的ID:
streamID.ms
,这是一个毫秒级别的时间戳,用于记录消息生成时的系统时间。streamID.seq
,相同时间戳内的消息序列号,如果同一毫秒内产生了多条消息,那么便会用这个序列号进行区分。每当这个毫秒时间戳被更新时,序列号会被清零重新计次。
Stream
typedef struct stream {rax *rax;uint64_t length;streamID last_id;rax *cgroups;} stream;
stream
这个数据结构便是Redis中消息队列的主体:
stream.rax
,这是一个基数树,以消息ID为Key,用来存储消息队列之中的消息。stream.length
,存储消息队列中消息的数量。stream.last_id
,记录上一次生成的消息ID,当插入新的消息时,会根据这个消息ID来生成新消息的消息ID。stream.cgroups
,同样是一个基数树,用于存储关联这个消息队列的消费者组,使用消费者组的名字作为基数树的Key。
在stream.rax
这个基数树中,Key为消息ID,而对应的Value则是以前面我们介绍过的紧凑列表listpack。简单描述一下消息队列之中基数树使用紧凑列表来存储消息的布局:
每条消息的消息内容都存储在一个紧凑列表之中。
每个紧凑列表中可以保存多条消息,每个消息都是这个紧凑列表的一个
entry
。紧凑列表会使用表中第一条消息的ID,记作
master_id
,作为这个紧凑列表在整个stream.rax
基数树中的Key。服务器在全局变量之中限制了单一紧凑列表的内存上限
redisServer.stream_node_max_bytes
以及单一紧凑列表最多可以容纳消息条数的上限redisServer.stream_node_max_entries
。当超过整个上限,Redis会新分配一个紧凑列表存储在基数树之中。当我们需要删除一条消息时,通常不是将这条消息从紧凑列表之中删除,而是将其设置为一个删除状态。因为删除会涉及到内存的移动已经重新分配,这样会大大地降低系统的性能。
下面我们可以来看一下在紧凑列表中,消息是按照何种格式存储的。
+-----------+-------+-------+-----+-------+|<MSGHeader>|<MSG-1>|<MSG-2>|.....|<MSG-N>|+-----------+-------+-------+-----+-------+
在MSGHeader
之中记录了当前这个紧凑列表之中的消息表头,后面的MSG-X
为存储这个紧凑列表之中的消息内容。
首先我们来看一下紧凑列表之中这个消息表头的内存分布:
+-------+---------+------------+---------+--/--+---------+---------+-+| count | deleted | num-fields | field_1 | field_2 | ... | field_N |0|+-------+---------+------------+---------+--/--+---------+---------+-+
count
之中记录这个紧凑列表之中存储的消息的数量。deleted
中记录了这个紧凑列表之中被标记为删除的消息的数量。fileds
这一组相关的数据是消息的field
数据,使用的是第一个加入这个紧凑列表之中的消息的field
数据。我们知道Redis消息队列之中的消息内容都是都按照field-string
这样格式的键值对,而能够加入同一个消息队列中的消息,其键值对中的fields
域大致应该是相同的。那么我们使用第一个加入这个紧凑列表中的消息的fields
域来作为这个列表的公共fields
域,如果后续消息的fields
域与第一个消息相同,那么可以公用这个紧凑列表中MSGHeader
里的fields
域,以达到节省内存的目的。num-fields
用于记录这个公用消息fields
域中field
的个数。field_XXX
则是用于记录每个field
的具体内容。最后一个
0
字符,相当于一个特殊的标记,用于标记MSGHeader
的结束。
在这个MSGHeader
之后,便是每条消息的具体内容,由于紧凑列表之中的消息有可能会被标记为已经删除的状态,因此Redis出了一个状态定义STREAM_ITEM_FLAG_DELETED
用于表示这个节点是否已经被删除。同时Redis还给出了一个状态定义STREAM_ITEM_FLAG_SAMEFIELDS
用于标记这个消息的fields
域是否与整个紧凑列表的fields
域相同,根据是否相同Redis对消息的内存分布格式给出两种定义。
首先我们看一下fields
与紧凑列表公共fields
域不同的消息的存储格式:
+-----+--------+----------+-------+-------+-/-+-------+-------+--------+|flags|entry-id|num-fields|field-1|value-1|...|field-N|value-N|lp-count|+-----+--------+----------+-------+-------+-/-+-------+-------+--------+
在上面这个内存分布之中:
flags
,这个字段以掩码的形式记录了当前消息的状态,前面提到的STREAM_ITEM_FLAG_DELETED
以及STREAM_ITEM_FLAG_SAMEFIELDS
都会被记录在这里。entry_id
,前面我们介绍了,在stream.rax
这个基数树中,Value是存储消息的紧凑列表,而Key是紧凑列表之中第一条消息的ID,我们称之为master_id
;同时我们也知道每条消息都用一个唯一的属于自己的ID。这里entry_id
便是用于记录消息的ID,只不过存储的不是消息的完整ID,而是记录的自身ID与master_id
的差值。num-fields
,由于这条消息的fields
域与公共fields
不相同,故此这里记录了该消息的键值对的个数。field-value
,这一系列字段记录了消息的键值对之中的内容。lp-count
,这个字段记录前面介绍的flags
,entry-id
等字段的个数,用于方便从后向前反向地进行遍历。
最后我们来看一下fields
域与公共fields
字段的
+-----+--------+-------+-/-+-------+--------+|flags|entry-id|value-1|...|value-N|lp-count|+-----+--------+-------+-/-+-------+--------+
这里类似flags
、entry-id
以及lp-count
字段与前面的含义相同,只不过由于省略了fields
域数据的存储,仅以value-XXX
来存储每个field
对应的value
的数据。
消费者组
typedef struct streamCG {streamID last_id;rax *pel;rax *consumers;} streamCG;
Redis使用streamCG
来表示一个消费者组:
streamCG.last_id
,消息派发游标,用于记录下一条需要发送给消费者的消息ID,每次消费者请求消息,消费者组便会将streamCG.last_id
对应的消息发送给消费者,并将streamCG.last_id
向后移动一位。streamCG.pel
,这是一个基数树结构,用于存储已经被发送但是未被确认消息的列表。streamCG.consumers
,同样是一个基数树结构,用于存储这个消费者组内的消费者,基数树的Key为消费者的名字,对应的Value
为后面介绍的streamConsumer
数据结构。
消费者
typedef struct streamConsumer {mstime_t seen_time;sds name;rax *pel;} streamConsumer;
streamConsumer
这个数据结构就是Redis用于表示消费者:
streamConsumer.seen_time
,记录消费者上次活动的时间戳。streamConsumer.name
,存储消费者的名字。streamConsumer.pel
,这个消费者对应的还没有确认的消息的列表。
未被确认的消息
typedef struct streamNACK {mstime_t delivery_time;uint64_t delivery_count;streamConsumer *consumer;} streamNACK;
streamNACK
这个数据结构就是表示已被派发但是没有被确认的消息:
streamNACK.delivery_time
,记录该条消息上次被派发的时间戳。streamNACK.delivery_count
,存储该条消息被派发的次数。streamNACK.consumer
,记录该条消息上次被派发的消费者。
Stream迭代器
typedef struct streamIterator {stream *stream;streamID master_id;uint64_t master_fields_count;unsigned char *master_fields_start;unsigned char *master_fields_ptr;int entry_flags;int rev;uint64_t start_key[2];uint64_t end_key[2];raxIterator ri;unsigned char *lp;unsigned char *lp_ele;unsigned char *lp_flags;unsigned char field_buf[LP_INTBUF_SIZE];unsigned char value_buf[LP_INTBUF_SIZE];} streamIterator;
在这个迭代器数据结构之中:
streamIterator.stream
,记录了当前迭代器关联的消息队列stream
的指针。streamIterator.master_id
,记录迭代器当前所在的紧凑列表对应的master_id
。streamIterator.master_fields_count
,记录迭代器当前所在的紧凑列表之中公共fields
域中field
的个数。streamIterator.master_fields_start
,记录迭代器当前所在的紧凑列表之中公共fields
域中第一个field
的指针。streamIterator.master_fields_ptr
,迭代器用于遍历当前在的紧凑列表中公共fields
域的指针。streamIterator.entry_flags
,迭代器当前遍历消息的标记字段。streamIterator.rev
,用于记录当前的迭代器是正向迭代器还是逆向迭代器。streamIterator.start_key
,记录这个迭代器初始化时的起始消息ID,默认是0-0
。streamIterator.end_key
,记录这个迭代器初始化时的终止消息ID,默认是UINT64_MAX-UINT64_MAX
。streamIterator.ri
,表示这个迭代器用于迭代遍历stream.rax
这个基数树的底层迭代器raxIterator
。streamIterator.lp
,记录迭代器当前指向的紧凑列表指针。streamIterator.lp_ele
,这是迭代器在遍历当前紧凑列表时使用的,指向列表元素的一个内部指针。streamIterator.lp_flags
,迭代器在当前的紧凑列表中指向列表中记录消息标记字段flags
内存的指针,streamIterator.entry_flags
可以通过这个字段解析获得。streamIterator.field_buf
,迭代器用于遍历某一条消息的键值对时,存储当前所看到的field
的内容。streamIterator.value_buf
,迭代器用于遍历某一条消息的键值对时,存储当前所看到的value
的内容。




