(二)kafka的架构中的主要角色
A. topic
定义:producer倾倒数据时最宏观的目标对象即为topic,一个topic往往会有多个partition分布在不同broker上。
组成:一个topic被创建时会被指定具体有多少个partition,producer通过算法实现partition在kafka cluster的分布均衡。
场景:consumer订阅时会明确指定订阅的topic,一般一个consumer group对应一个topic,一个consumer对应一个partition时最高效。
B. partition
组成:producer实际倾倒数据时是指向broker的某个partition目录的,partition目录下是一个个segment文件来存储实际信息的。
场景:默认情况下分区数为3,若topic名称为topicA,此时拥有partition为:topicA-0,topicA-1 & topicA-2。若kafka集群仅有3个broker组成,那么此时分区的leader副本能均匀分布在3个broker上。
-
合理计算partition数量
- 总体上,partition数量一般由消费速度来决定。
- 若单个consumer的消费速度为50Mb/s,该topic的producer的总写入速度是1Gb/s,则可计算出消费该topic需要20个consumers,因此可规划该topic的partition数目为20个。
-
partition大小的重要性
过多的partition对应产生有大量的目录,虽然单个目录下的文件是顺序写入,但宏观看大量目录切换破坏了追加写入的初衷。
故障转移时controller是单线程处理的,若有大量分区,那么故障发生时这些leader副本的partition都暂时无法提供读写服务,会导致整个故障转移很慢。当前优化了zk的批量异步写、批量计算和日志优化来缓解该问题。
大量partition就会导致metadata过大,集群状态若变化了更新metadata的过程也会太沉重。
broker会于后台定时清理kafka日志文件,分区过大导致文件数量过多,线程因此会处理压力较大,导致磁盘存在过载风险。
C. replica
1. 基础概述
- 定义:为保证高可用性,kafka的每个partition都有自己的副本replica,被称为follower副本。follower副本是由其他broker向leader副本所在broker主动拉取数据的(pull模式)。
- 场景:延续上方partition中的场景,若topicA的replication factor为2,那么topicA-0的leader副本置于broker0上,对应的follower副本放在broker1上。
2. 副本下的ISR
- ISR定义:全称in-sync-replicas,指与leader保持同步的replicas。AR是所有分配的副本,OSR是未同步的副本,AR=ISR+OSR。
-
同步基础
- ISR是动态变化的,通过ISR是否消息都同步,确认我们写入消息是否成功。
- 结合参数acks,若acks=-1时需ISR中的副本数量为
min.insync.replicas时被认为写入成功,若acks=1则仅需所有ISR中的副本数量有1个即可。
- 场景:为确保有两个replicas同步数据才认定写入成功,一般会设置
min.insync.replicas=2并搭配acks=-1。若min.insync.replicas=1,那么acks=1 OR -1的实际效力是相同的。
-
伸缩机制
-
replica.lag.max.messages默认值是4000,即follower落后leader4000个消息后就被踢出ISR。 - 若在高吞吐场景,4k延迟会导致ISR原数据频繁更新,kafka controller替代zk watcher收集信息依然会压力偏大;若是低吞吐场景,4k延迟不够灵敏,明明延迟了却未及时反馈。
- 现在使用
replica.lag.time.max.ms=10000作为默认参数,即10s内follower无法追上leader的LEO(log end offset)日志末端位移,即被判定踢出ISR。
-
-
LEO & HW
- LEO是日志末端位移,相对的是LBO(log base offset)日志初始位移。
- HW(high watermark)则是副本的高水位,ISR集合中的min LEO决定了HW,HW是消费者实际能消费到的位置。
-
isr-expiration参数
- 每台broker有定时任务检查本台leader的isr列表有无失效副本,周期是
replica.lag.time.max.ms/2,该参数默认值是10s,所以周期为5s。 -
lastCaughtUpTimeMs在follwer的LEO和leader的LEO相等时更新,根据当前时间now减去follower的lastCaughtUpTimeMs,若该值大于
replica.lag.time.max.ms,则replicas失效。
- 每台broker有定时任务检查本台leader的isr列表有无失效副本,周期是
-
失效副本的修复机制
- 若失效的follower副本的LEO再度等于leader副本的HW时,可能会被拉回到ISR中。
- 为防止起变化过频繁,Kafka要求必须在ISR发生变化超过5s,且距离上次元数据写入zk超过60s时,才能重新加入ISR中。
D. controller
1. 基础介绍
概述:用于服务端容灾和集群管理,kafka的controller并非主从结构,其主要负责metadata收集以降低zk watch机制的负载,由于kafka partition太多容易导致zk压力较大。
-
故障转移顺序:整体流程是单线程顺序执行的,在执行过程中故障发生的节点下所有leader副本都不可提供读写,对业务有较大的负面影响。controller本身是依赖于zk的临时节点,会注册对应的监听器,若出现事情就触发对应的handler。针对感知到故障发生时,controller会通知所有broker将有必要的follower副本变为leader副本,具体流程的执行步骤如下
- 计算新leader
- 更新zk数据
- 通知其他节点将其拥有的follower副本转变成leader
- 更新元数据
重要性:在实现故障转移功能外,controller的存在也降低了kafka对于zk的依赖,使别的broker只需要注册少量watcher,仅controller本身仍需注册大量watcher。
2. controller工作流程详解
集群状态同步:controller负责管理broker、topic & isr等,若感知到集群状态有变化,则将新的metadata同步到其他broker节点并更新集群的metadata。
-
选举controller的顺序
初次选举:broker启动后便会去zk上注册临时节点 - /brokers/ids/brokerId & /controller,在前者写入自己的brokerId & timestamp。一般情况下,基于节点不存在就能创建成功的原则,谁先启动就会先抢到controller的角色。若未抢到该角色(即controller临时节点已存在),则会获取到当选controller的brokerId,并创建一个watcher,一旦controller节点被删除,所有注册watcher的节点就会尝试抢占;
宕机后选举:若原controller宕机,注册的zk临时节点就消失(被删除)了,此时会触发其他broker的zk临时节点所设置的监听事件,触发竞选;
确保唯一性:每次更新完controller后,zk的持久化节点 - /controller_epoch就会更新(+1),保证全局仅一个controller。