暂无图片
暂无图片
暂无图片
暂无图片
暂无图片

Redis源码学习(59)-Redis可持久化的消息队列(3)

马基雅维利incoding 2021-03-02
249

Redis中Stream的基础数据结构

消息ID

    typedef struct streamID {
    uint64_t ms;
    uint64_t seq;
    } streamID;

    Redis使用streamID
    这个数据结构来表示消息的ID:

    1. streamID.ms
      ,这是一个毫秒级别的时间戳,用于记录消息生成时的系统时间。

    2. streamID.seq
      ,相同时间戳内的消息序列号,如果同一毫秒内产生了多条消息,那么便会用这个序列号进行区分。每当这个毫秒时间戳被更新时,序列号会被清零重新计次。

    Stream

      typedef struct stream {
      rax *rax;
      uint64_t length;
      streamID last_id;
      rax *cgroups;
      } stream;

      stream
      这个数据结构便是Redis中消息队列的主体:

      1. stream.rax
        ,这是一个基数树,以消息ID为Key,用来存储消息队列之中的消息。

      2. stream.length
        ,存储消息队列中消息的数量。

      3. stream.last_id
        ,记录上一次生成的消息ID,当插入新的消息时,会根据这个消息ID来生成新消息的消息ID。

      4. stream.cgroups
        ,同样是一个基数树,用于存储关联这个消息队列的消费者组,使用消费者组的名字作为基数树的Key

      stream.rax
      这个基数树中,Key为消息ID,而对应的Value则是以前面我们介绍过的紧凑列表listpack。简单描述一下消息队列之中基数树使用紧凑列表来存储消息的布局:

      1. 每条消息的消息内容都存储在一个紧凑列表之中。

      2. 每个紧凑列表中可以保存多条消息,每个消息都是这个紧凑列表的一个entry

      3. 紧凑列表会使用表中第一条消息的ID,记作master_id
        ,作为这个紧凑列表在整个stream.rax
        基数树中的Key

      4. 服务器在全局变量之中限制了单一紧凑列表的内存上限redisServer.stream_node_max_bytes
        以及单一紧凑列表最多可以容纳消息条数的上限redisServer.stream_node_max_entries
        。当超过整个上限,Redis会新分配一个紧凑列表存储在基数树之中。

      5. 当我们需要删除一条消息时,通常不是将这条消息从紧凑列表之中删除,而是将其设置为一个删除状态。因为删除会涉及到内存的移动已经重新分配,这样会大大地降低系统的性能。

      下面我们可以来看一下在紧凑列表中,消息是按照何种格式存储的。

      +-----------+-------+-------+-----+-------+|<MSGHeader>|<MSG-1>|<MSG-2>|.....|<MSG-N>|+-----------+-------+-------+-----+-------+

      MSGHeader
      之中记录了当前这个紧凑列表之中的消息表头,后面的MSG-X
      为存储这个紧凑列表之中的消息内容。

      首先我们来看一下紧凑列表之中这个消息表头的内存分布:

      +-------+---------+------------+---------+--/--+---------+---------+-+| count | deleted | num-fields | field_1 | field_2 | ... | field_N |0|+-------+---------+------------+---------+--/--+---------+---------+-+
      1. count
        之中记录这个紧凑列表之中存储的消息的数量。

      2. deleted
        中记录了这个紧凑列表之中被标记为删除的消息的数量。

      3. fileds
        这一组相关的数据是消息的field
        数据,使用的是第一个加入这个紧凑列表之中的消息的field
        数据。我们知道Redis消息队列之中的消息内容都是都按照field-string
        这样格式的键值对,而能够加入同一个消息队列中的消息,其键值对中的fields
        域大致应该是相同的。那么我们使用第一个加入这个紧凑列表中的消息的fields
        域来作为这个列表的公共fields
        域,如果后续消息的fields
        域与第一个消息相同,那么可以公用这个紧凑列表中MSGHeader
        里的fields
        域,以达到节省内存的目的。

        1. num-fields
          用于记录这个公用消息fields
          域中field
          的个数。

        2. field_XXX
          则是用于记录每个field
          的具体内容。

      4. 最后一个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|+-----+--------+----------+-------+-------+-/-+-------+-------+--------+

      在上面这个内存分布之中:

      1. flags
        ,这个字段以掩码的形式记录了当前消息的状态,前面提到的STREAM_ITEM_FLAG_DELETED
        以及STREAM_ITEM_FLAG_SAMEFIELDS
        都会被记录在这里。

      2. entry_id
        ,前面我们介绍了,在stream.rax
        这个基数树中,Value是存储消息的紧凑列表,而Key是紧凑列表之中第一条消息的ID,我们称之为master_id
        ;同时我们也知道每条消息都用一个唯一的属于自己的ID。这里entry_id
        便是用于记录消息的ID,只不过存储的不是消息的完整ID,而是记录的自身ID与master_id
        的差值。

      3. num-fields
        ,由于这条消息的fields
        域与公共fields
        不相同,故此这里记录了该消息的键值对的个数。

      4. field-value
        ,这一系列字段记录了消息的键值对之中的内容。

      5. 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
        来表示一个消费者组

        1. streamCG.last_id
          ,消息派发游标,用于记录下一条需要发送给消费者的消息ID,每次消费者请求消息,消费者组便会将streamCG.last_id
          对应的消息发送给消费者,并将streamCG.last_id
          向后移动一位。

        2. streamCG.pel
          ,这是一个基数树结构,用于存储已经被发送但是未被确认消息的列表。

        3. streamCG.consumers
          ,同样是一个基数树结构,用于存储这个消费者组内的消费者,基数树的Key消费者的名字,对应的Value
          为后面介绍的streamConsumer
          数据结构。

        消费者

          typedef struct streamConsumer {
          mstime_t seen_time;
          sds name;
          rax *pel;
          } streamConsumer;

          streamConsumer
          这个数据结构就是Redis用于表示消费者

          1. streamConsumer.seen_time
            ,记录消费者上次活动的时间戳。

          2. streamConsumer.name
            ,存储消费者的名字。

          3. streamConsumer.pel
            ,这个消费者对应的还没有确认的消息的列表。

          未被确认的消息

            typedef struct streamNACK {
            mstime_t delivery_time;
            uint64_t delivery_count;
            streamConsumer *consumer;
            } streamNACK;

            streamNACK
            这个数据结构就是表示已被派发但是没有被确认的消息:

            1. streamNACK.delivery_time
              ,记录该条消息上次被派发的时间戳。

            2. streamNACK.delivery_count
              ,存储该条消息被派发的次数。

            3. 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;

              在这个迭代器数据结构之中:

              1. streamIterator.stream
                ,记录了当前迭代器关联的消息队列stream
                的指针。

              2. streamIterator.master_id
                ,记录迭代器当前所在的紧凑列表对应的master_id

              3. streamIterator.master_fields_count
                ,记录迭代器当前所在的紧凑列表之中公共fields
                域中field
                的个数。

              4. streamIterator.master_fields_start
                ,记录迭代器当前所在的紧凑列表之中公共fields
                域中第一个field
                的指针。

              5. streamIterator.master_fields_ptr
                ,迭代器用于遍历当前在的紧凑列表中公共fields
                域的指针。

              6. streamIterator.entry_flags
                ,迭代器当前遍历消息的标记字段。

              7. streamIterator.rev
                ,用于记录当前的迭代器是正向迭代器还是逆向迭代器。

              8. streamIterator.start_key
                ,记录这个迭代器初始化时的起始消息ID,默认是0-0

              9. streamIterator.end_key
                ,记录这个迭代器初始化时的终止消息ID,默认是UINT64_MAX-UINT64_MAX

              10. streamIterator.ri
                ,表示这个迭代器用于迭代遍历stream.rax
                这个基数树的底层迭代器raxIterator

              11. streamIterator.lp
                ,记录迭代器当前指向的紧凑列表指针。

              12. streamIterator.lp_ele
                ,这是迭代器在遍历当前紧凑列表时使用的,指向列表元素的一个内部指针。

              13. streamIterator.lp_flags
                ,迭代器在当前的紧凑列表中指向列表中记录消息标记字段flags
                内存的指针,streamIterator.entry_flags
                可以通过这个字段解析获得。

              14. streamIterator.field_buf
                ,迭代器用于遍历某一条消息的键值对时,存储当前所看到的field
                的内容。

              15. streamIterator.value_buf
                ,迭代器用于遍历某一条消息的键值对时,存储当前所看到的value
                的内容。

              文章转载自马基雅维利incoding,如果涉嫌侵权,请发送邮件至:contact@modb.pro进行举报,并提供相关证据,一经查实,墨天轮将立刻删除相关内容。

              评论