kafka学习小结
《胡夕.ApacheKafka实战(Kindle位置2363-2369).电子工业出版社.Kindle版本.》
Kafka的高性能的原因
https://zhuanlan.zhihu.com/p/96957920
大量使用操作系统页缓存,内存操作速度快命中率高;
不直接参与io,而是交给操作系统做;
零拷贝减少不必要的复制;
采用追加写提高读写速度;
Kafka消息的格式
使用二进制数组而不是对象形式,紧凑节约空间,包括key(用于分区判断),value,时间戳,其他信息(透明)
如何保证高可用
分片、副本、leader、follower,zk在整个集群中选举出一个Broker作为Controller,Controller为所有Topic的所有Partition指定Leader及Follower,减轻Zookeeper负载
设计架构
producer:通过分区器向缓冲区发送数据,有个线程专门用于发送数据
consumer:心跳线程+主线程,一个主线程轮流对接多个partition进行消费,位移由kafka自行保管
broker:分区+副本
controller,是zk选出的某个broker,创建、删除主题,增加分区并分配leader分区、集群Broker管理(新增 Broker、Broker 主动关闭、Broker 故障)、preferred leader选举、分区重分配
coordinator,每个KafkaServer都有一个GroupCoordinator实例,管理多个消费者组,主要用于offset位移管理和ConsumerRebalance。consumergroup在执行rebalance之前必须首先确定coordinator所在的broker,并创建与该broker相互通信的Socket连接。确定coordinator的算法与确定offset被提交到__consumer_offsets目标分区的算法是相同的。算法如下。?计算Math.abs(groupID.hashCode)%offsets.topic.num.partitions参数值(默认是50),假设是10。?寻找__consumer_offsets分区10的leader副本所在的broker,该broker即为这个group的coordinator。成功连接coordinator之后便可以执行rebalance操作。目前rebalance主要分为两步:加入组和同步更新分配方案。
集群的副本同步机制
leader接收数据、副本同步数据
如果leader crash时,ISR为空怎么办?
kafka在Broker端提供了一个配置参数:unclean.leader.election,这个参数有两个值:
true(默认):允许不同步副本成为leader,由于不同步副本的消息较为滞后,此时成为leader,可能会出现消息不一致的情况。
false:不允许不同步副本成为leader,此时如果发生ISR列表为空,会一直等待旧leader恢复,降低了可用性。
Kafka的生产者发送消息流程是怎么样

如何保证消息不丢失
producer端配置
acks=all 确保所有的follower副本都更新
retries=Integer.MAX_VALUE 一直重试
max.in.flight.requests.per.connection=1 一直只允许发送一条消息,避免乱序
Callback的失败处理逻辑中显式调用KafkaProducer.close(0)
broker端配置
unclean.leader.election.enable=false 关闭unclean leader选举, 即不允许非ISR中的副本被选举为leader, 从而避免broker端因日志水位截断而造成的消息丢失。
replication.factor>=3
min.insync.replicas>1 用于控制某条消息至少被写入到ISR中的多少个副本才算成功
replication.factor=min.insyn.replicas+1
如何增加写的性能
acks=1
信息压缩
适当提高batch.size,当batch满了的时候,producer会发送batch中的所有消息。
适当提高linger.ms,linger.ms参数就是控制消息发送延时行为的。该参数默认值是0,表示消息需要被立即发送,无须关心batch是否已被填满,大多数情况下这是合理的,毕竟我们总是希望消息被尽可能快地发送。不过这样做会拉低producer吞吐量,毕竟producer发送的每次请求中包含的消息数越多,producer就越能将发送请求的开销摊薄到更多的消息上,从而提升吞吐量。
多线程
判断broker异常
故障转移通常是以“心跳”或“会话”的机制来实现的,即只要主服务器与备份服务器之间的心跳无法维持或主服务器注册到服务中心的会话超时过期了,那么就认为主服务器已无法正常运行,集群会自动启动某个备份服务器来替代主服务器的工作。Kafka服务器支持故障转移的方式就是使用会话机制。每台Kafka服务器启动后会以会话的形式把自己注册到ZooKeeper服务器上。一旦该服务器运转出现问题,与ZooKeeper的会话便不能维持从而超时失效,此时Kafka集群会选举出另一台服务器来完全代替这台服务器继续提供服务
Kafka怎么实现分区策略,怎么实现负载均衡
Kafkaproducer提供了一个默认的分区器。对于每条待发送的消息而言,如果该消息指定了key,那么该partitioner会根据key的哈希值来选择目标分区;若这条消息没有指定key,则partitioner使用轮询的方式确认目标分区——这样可以最大限度地确保消息在所有分区上的均匀性。当然producer的API赋予了用户自行指定目标分区的权力,即用户可以在消息发送时跳过partitioner直接指定要发送到的分区。另外,producer也允许用户实现自定义的分区策略而非使用默认的partitioner,这样用户可以很灵活地根据自身的业务需求确定不同的分区策略。
如果kafka消费者消费超时会发生什么?怎么避免kafka的消费超时
组rebalance触发的条件有以下3个。?组成员发生变更,比如新consumer加入组,或已有consumer主动离开组,再或是已有consumer崩溃时则触发rebalance。?组订阅topic数发生变更,比如使用基于正则表达式的订阅,当匹配正则表达式的新topic被创建时则会触发rebalance。?组订阅topic的分区数发生变更,比如使用命令行脚本增加了订阅topic的分区数。
真实应用场景中引发rebalance最常见的原因就是违背了第一个条件,特别是consumer崩溃的情况。这里的崩溃不一定就是指consumer进程“挂掉”或consumer进程所在的机器宕机。当consumer无法在指定的时间内完成消息的处理,那么coordinator就认为该consumer已经崩溃,从而引发新一轮rebalance。
(1)参数调整session.timeout.ms、max.poll.records和max.poll.interval.ms
session.timeout.ms参数被明确为“coordinator检测失败的时间”。因此在实际使用中,用户可以为该参数设置一个比较小的值让coordinator能够更快地检测consumer崩溃的情况,从而更快地开启rebalance,避免造成更大的消费滞后(consumerlag)。目前该参数的默认值是10秒。
max.poll.interval.ms。在一个典型的consumer使用场景中,用户对于消息的处理可能需要花费很长时间。这个参数就是用于设置消息处理逻辑的最大时间的。
max.poll.records该参数控制单次poll调用返回的最大消息数。比较极端的做法是设置该参数为1,那么每次poll只会返回1条消息。如果用户发现consumer端的瓶颈在poll速度太慢,可以适当地增加该参数的值。如果用户的消息处理逻辑很轻量,默认的500条消息通常不能满足实际的消息处理速度。
(2)多worker线程消费
rebalance分配策略
rebalance时group下所有的consumer都会协调在一起共同参与分区分配,这是如何完成的呢?Kafka新版本consumer默认提供了3种分配策略,分别是range策略、round-robin策略和sticky策略。
所谓的分配策略决定了订阅topic的每个分区会被分配给哪个consumer。range策略主要是基于范围的思想。它将单个topic的所有分区按照顺序排列,然后把这些分区划分成固定大小的分区段并依次分配给每个consumer;round-robin策略则会把所有topic的所有分区顺序摆开,然后轮询式地分配给各个consumer。最新发布的sticky策略有效地避免了上述两种策略完全无视历史分配方案的缺陷,采用了“有黏性”的策略对所有consumer实例进行分配,可以规避极端情况下的数据倾斜并且在两次rebalance间最大限度地维持了之前的分配方案。
通常意义上认为,如果group下所有consumer实例的订阅是相同,那么使用round-robin会带来更公平的分配方案,否则使用range策略的效果更好。新版本consumer默认的分配策略是range。
rebalance过程
consumergroup在执行rebalance之前必须首先确定coordinator所在的broker,并创建与该broker相互通信的Socket连接。确定coordinator的算法与确定offset被提交到__consumer_offsets目标分区的算法是相同的。算法如下。?计算Math.abs(groupID.hashCode)%offsets.topic.num.partitions参数值(默认是50),假设是10。?寻找__consumer_offsets分区10的leader副本所在的broker,该broker即为这个group的coordinator。成功连接coordinator之后便可以执行rebalance操作。目前rebalance主要分为两步:加入组和同步更新分配方案。?加入组:这一步中组内所有consumer(即group.id相同的所有consumer实例)向coordinator发送JoinGroup请求。当收集全JoinGroup请求后,coordinator从中选择一个consumer担任group的leader,并把所有成员信息以及它们的订阅信息发送给leader。特别需要注意的是,group的leader和coordinator不是一个概念。leader
是某个consumer实例,coordinator通常是Kafka集群中的一个broker。另外leader而非coordinator负责为整个group的所有成员制定分配方案。?同步更新分配方案:这一步中leader开始制定分配方案,即根据前面提到的分配策略决定每个consumer都负责哪些topic的哪些分区。一旦分配完成,leader会把这个分配方案封装进SyncGroup请求并发送给coordinator。比较有意思的是,组内所有成员都会发送SyncGroup请求,不过只有leader发送的SyncGroup请求中包含了分配方案。coordinator接收到分配方案后把属于每个consumer的方案单独抽取出来作为SyncGroup请求的response返还给各自的consumer。
选举哪些地方
zk选broker、broker选副本leader、coordinator选consumer的leader
如何设计保证kakfa中有某个特征的消息的是严格按照顺序消费的
设置kafka只有一个分区,或者producer指定一个partition发送或者指定相同的key将消息打到同一个partition中,因为kafka只保证同一个分区的有序性,不保证全局的有序性。
在保证使用一个分区的基础上,需要设置以下两种参数保证有序:(1)设置retries>0,并且max.in.flight.requests.per.connection=1,前者保证数据在临时异常时可以再次发送,后者保证在缓冲区中只有一个消息待发送,不会因为重试造成乱序;(2)enable.idempotence=true时可以使max.in.flight.requests.per.connection>1,因为幂等性可以保证无法跨越一个消息发送,也可以保证重复消息被过滤。
幂等性说明,对于每个PID,该Producer发送消息的每个
注意,在多分区的情况下,需要使用独立consumer直接利用assign方法订阅指定的分区。
Kafka如何实现消息幂
幂等性producer,瞬时的发送错误可能导致producer端出现重试,同一条消息被producer发送多次,但在broker端这条消息只会被写入日志一次。显式地设置producer端的新参数enable.idempotence为true。
幂等性producer的设计思路类似于TCP的工作方式。发送到broker端的每批消息都会被赋予一个序列号(sequencenumber)用于消息去重。但是和TCP不同的是,这个序列号不会被丢弃,相反Kafka会把它们保存在底层日志中,这样即使分区的leader副本挂掉,新选出来的leaderbroker也能执行消息去重工作。除了序列号,Kafka还会为每个producer实例分配一个producerid(下称PID)。producer在初始化时必须分配一个PID。PID分配的过程对用户来说是完全透明的,因此不会为用户所见。消息要被发送到的每个分区都有对应的序列号值,它们总是从0开始并且严格单调增加。对于PID、分区和序列号的关系,用户可以设想一个Map,key就是(PID,分区号),value就是序列号。即每对(PID,分区号)都有对应的序列号值。若发送消息的序列号小于或等于broker端保存的序列号,那么broker会拒绝这条消息的写入操作。这种设计确保了即使出现重试操作,每条消息也只会被保存在日志中一次。不过,由于每个新的producer实例都会被分配不同的PID,当前设计只能保证单个producer实例的EOS语义,而无法实现多个producer实例一起提供EOS语义。这一点要特别注意。
对于接收的每条消息,如果其序号比Broker维护的序号)大一,则Broker会接受它,否则将其丢弃:如果消息序号比Broker维护的序号差值比一大,说明中间有数据尚未写入,即乱序,此时Broker拒绝该消息,Producer抛出InvalidSequenceNumber如果消息序号小于等于Broker维护的序号,说明该消息已被保存,即为重复消息,Broker直接丢弃该消息,Producer抛出DuplicateSequenceNumberSender发送失败后会重试,这样可以保证每个消息都被发送到broker
如何保证消息只消费一次
幂等性和事务