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

KAFKA的消费时OFFSET管理

袖手安坐 2020-07-04
435

基本概念介绍

  • topic

    消息队列

  • partition

    topic下的存储单元,类似分片(sharding),每个partition对应一个实体文件

  • offset

    partition内每条消息的ID,单个partition内递增

  • broker

    kafka集群的节点,负责存储消息,也是生产消费消息的接入点

  • consumer_group

    消费组。同一消费组对应的同一topic的下同一个partition只有一个当前消费的offset点的记录。也就是说kafka记录哪些消息已经被消费用的是consumer_group, topic, partiton
    这个三元组为key,value为offset

  • member

    同一consumer_group内的多个消费者




消费流程

  1. 加入消费组

  2. 找到消费组的协调者 group coordinator (某个broker)

    https://kafka.apache.org/protocol#The_Messages_FindCoordinator

      FindCoordinator Request (Version: 0) => key
      key => STRING

      向 group coordinator 发起join group请求

      https://kafka.apache.org/protocol#The_Messages_JoinGroup

        JoinGroup Response (Version: 0) => error_code generation_id protocol_name leader member_id [members]
        error_code => INT16
        generation_id => INT32
        protocol_name => STRING
        leader => STRING
        member_id => STRING
        members => member_id metadata
        member_id => STRING
        metadata => BYTES
          JoinGroup Request (Version: 0) => group_id session_timeout_ms member_id protocol_type [protocols]
          group_id => STRING
          session_timeout_ms => INT32
          member_id => STRING
          protocol_type => STRING
          protocols => name metadata
          name => STRING
          metadata => BYTES

          group coordinator 决定 group leader

          group leader拿到member列表

          group leader决定 partition-member
          对应关系


          获取partition和broker的对应关系,一般sdk会做缓存,然后定时更新meta信息


          根据topic的信息,默认使用 RoundRobinAssignor
          算法得出partition-member
          对应关系

          如果meta信息更新,重新计算

            Metadata Request (Version: 0) => [topics]
            topics => name
            name => STRING
              Metadata Response (Version: 0) => [brokers] [topics]
              brokers => node_id host port
              node_id => INT32
              host => STRING
              port => INT32
              topics => error_code name [partitions]
              error_code => INT16
              name => STRING
              partitions => error_code partition_index leader_id [replica_nodes] [isr_nodes]
              error_code => INT16
              partition_index => INT32
              leader_id => INT32
              replica_nodes => INT32
              isr_nodes => INT32

              group leader发送partition-member
              对应关系给  group coordinator

              https://kafka.apache.org/protocol#The_Messages_SyncGroup



                SyncGroup Response (Version: 0) => error_code assignment
                error_code => INT16
                assignment => BYTES

                group coordinator发送partition-member
                对应关系给所有的group member



                所有的group member都发起SyncGroup
                的请求,得到请求分配信息,只有group leader发起请求时

                  SyncGroup Response (Version: 1) => throttle_time_ms error_code assignment
                  throttle_time_ms => INT32
                  error_code => INT16
                  assignment => BYTES
                • 消费

                  • 发送心跳,检测是否需要rebalance,更新partition-member
                    对应关系

                    HeartbeatRequest => GroupId GenerationId MemberId
                    GroupId => string
                    GenerationId => int32
                    MemberId => string

                  • 获取offset

                    https://kafka.apache.org/protocol#The_Messages_OffsetFetch

                    OffsetFetch Request (Version: 1) => group_id [topics]
                    group_id => STRING
                    topics => topic [partitions]
                    topic => STRING
                    partitions => partition
                    partition => INT32

                  • 根据offset从对应的broker获取消息数据

                    https://kafka.apache.org/protocol#The_Messages_Fetch

                    Fetch Request (Version: 0) => replica_id max_wait_time min_bytes [topics]
                    replica_id => INT32
                    max_wait_time => INT32
                    min_bytes => INT32
                    topics => topic [partitions]
                    topic => STRING
                    partitions => partition fetch_offset partition_max_bytes
                    partition => INT32
                    fetch_offset => INT64
                    partition_max_bytes => INT32

                • ACK提交

                  https://kafka.apache.org/protocol#The_Messages_OffsetCommit

                  OffsetCommit Request (Version: 0) => group_id [topics]
                  group_id => STRING
                  topics => name [partitions]
                  name => STRING
                  partitions => partition_index committed_offset committed_metadata
                  partition_index => INT32
                  committed_offset => INT64
                  committed_metadata => NULLABLE_STRING

                  常用的两种ACK方法

                  对于严谨的服务,正确消费消息然后调用API进行commit,如果worker异常,消息会被再次处理,排除单一consumer worker的异常。定期commit则有可能造成消息没有正常处理。

                  更好的ack处理方式应该包含未ack超时,最大重试次数和死信队列这三个功能,这也是消息队列完整ack的必要功能。这里主要关于Kafka,就不展开了。

                  • 调用API进行commit

                  • sdk定期commit





                关于原理的思考

                1. 队列的ack只能对partition纬度进行

                2. 一个partiton只能被一个worker消费

                  消费者的并发数 不会 大于 partition的数量

                  每个worker的消费能力不同怎么办?

                3. partition增加,或是worker变化,都需要改变member-partion的对应关系,触发reblance

                        一个异常节点的频繁加入退出consumer group会导致所有的消费异常














                消费相关的常见问题

                1. 消息的过期时间

                  两种过期方式,按时间和按单个partition大小,可以以topic为维度自定义。过期之后消息可以在broker级定义为删除或是其他处理策略,对于API层面就是消息访问不到了。

                  两种方式同时配置时,先触发过期清理的策略生效。

                  retention.bytes

                  retention.ms

                2. offset过期时间

                  consumer提交的offset是有过期时间的,过期之后取不到,和初次消费的处理方案一致。

                  offset.retention.minutes
                  broker级别配置,默认2.0版本前1天,2.0后版本默认7天。

                3. rebalance

                  同一costumer group的各个member获取member-partition对应关系的过程,以下情况触发rebalance过程



                  目前的rebalance过程会使得所有的worker停止工作,

                  增量的rebalance方式可能会出现在未来的版本中

                  • member 加入/离开 consumer group

                  • member 心跳超时 heartbeat.interval.ms

                    发给某个broker节点,group coordinator

                  • partition 数量变化

                4. group coordinator的选定方式

                  offsets.topic.replication.factor
                   参数控制

                5. 单条消息大小

                  topic相关max.message.bytes

                  producer相关 max.request.size

                6. 新增partition后,消费延迟consumer lag
                  变化

                  默认round robin方式生产,消费也是一个member对应一个partition,consumer lag
                  会缓慢变化。

                7. 用offset获取消息,获取不到时的处理策略

                  offset过期,或是初次消费 通用遵循处理策略

                  consumer auto.offset.reset


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

                评论