基本概念介绍
topic
消息队列
partition
topic下的存储单元,类似分片(sharding),每个partition对应一个实体文件
offset
partition内每条消息的ID,单个partition内递增
broker
kafka集群的节点,负责存储消息,也是生产消费消息的接入点
consumer_group
消费组。同一消费组对应的同一topic的下同一个partition只有一个当前消费的offset点的记录。也就是说kafka记录哪些消息已经被消费用的是
consumer_group, topic, partiton
这个三元组为key,value为offsetmember
同一consumer_group内的多个消费者
消费流程
加入消费组
消费
发送心跳,检测是否需要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 => INT32ACK提交
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
找到消费组的协调者 group coordinator (某个broker)
https://kafka.apache.org/protocol#The_Messages_FindCoordinator
FindCoordinator Request (Version: 0) => keykey => 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 => INT16generation_id => INT32protocol_name => STRINGleader => STRINGmember_id => STRINGmembers => member_id metadatamember_id => STRINGmetadata => BYTES
JoinGroup Request (Version: 0) => group_id session_timeout_ms member_id protocol_type [protocols]group_id => STRINGsession_timeout_ms => INT32member_id => STRINGprotocol_type => STRINGprotocols => name metadataname => STRINGmetadata => 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 => namename => STRING
Metadata Response (Version: 0) => [brokers] [topics]brokers => node_id host portnode_id => INT32host => STRINGport => INT32topics => error_code name [partitions]error_code => INT16name => STRINGpartitions => error_code partition_index leader_id [replica_nodes] [isr_nodes]error_code => INT16partition_index => INT32leader_id => INT32replica_nodes => INT32isr_nodes => INT32
group leader发送partition-member
对应关系给 group coordinator
https://kafka.apache.org/protocol#The_Messages_SyncGroup
SyncGroup Response (Version: 0) => error_code assignmenterror_code => INT16assignment => BYTES
group coordinator发送partition-member
对应关系给所有的group member
所有的group member都发起SyncGroup
的请求,得到请求分配信息,只有group leader发起请求时
SyncGroup Response (Version: 1) => throttle_time_ms error_code assignmentthrottle_time_ms => INT32error_code => INT16assignment => BYTES
关于原理的思考
队列的ack只能对partition纬度进行
一个partiton只能被一个worker消费
消费者的并发数 不会 大于 partition的数量
每个worker的消费能力不同怎么办?
partition增加,或是worker变化,都需要改变member-partion的对应关系,触发reblance
一个异常节点的频繁加入退出consumer group会导致所有的消费异常
消费相关的常见问题
消息的过期时间
两种过期方式,按时间和按单个partition大小,可以以topic为维度自定义。过期之后消息可以在broker级定义为删除或是其他处理策略,对于API层面就是消息访问不到了。
两种方式同时配置时,先触发过期清理的策略生效。
retention.bytesretention.msoffset过期时间
consumer提交的offset是有过期时间的,过期之后取不到,和初次消费的处理方案一致。
offset.retention.minutes
broker级别配置,默认2.0版本前1天,2.0后版本默认7天。rebalance
同一costumer group的各个member获取member-partition对应关系的过程,以下情况触发rebalance过程
目前的rebalance过程会使得所有的worker停止工作,
增量的rebalance方式可能会出现在未来的版本中

member 加入/离开 consumer group
member 心跳超时
heartbeat.interval.ms发给某个broker节点,group coordinator
partition 数量变化
group coordinator的选定方式
offsets.topic.replication.factor
参数控制单条消息大小
topic相关
max.message.bytesproducer相关
max.request.size新增partition后,消费延迟
consumer lag
变化默认round robin方式生产,消费也是一个member对应一个partition,
consumer lag
会缓慢变化。用offset获取消息,获取不到时的处理策略
offset过期,或是初次消费 通用遵循处理策略
consumer
auto.offset.reset




