Kafka 核心面试题
文章目录
1. Kafka 如何保证消息不丢失
Kafka 是一个用来实现异步消息通信的中间件,它的整个架构由Producer、Consumer、Broker组成。
所以 对 kafka 保证消息不丢失这个问题,可以从三个方面来考虑和实现:
- 首先是Producer端,需要确保消息能够到达Broker并实现消息存储,在这个层面,有可能出现网络问题,导致消息发送失败,所以,针对Producer端,可以通过2种方式来避免消息丢失:
Producer默认是异步发送消息,这种情况下要确保消息发送成功,有两个方法
a. 把异步发送改成同步发送,这样producer就能实时知道消息发送的结果。
b. 添加异步回调函数来监听消息发送的结果,如果发送失败,可以在回调中重试。 - 然后是Broker端,Broker需要确保Producer发送过来的消息不会丢失,也就是只需要把消息持久化到磁盘就可以了。
- 但是,Kafka为了提升性能,采用了异步批量刷盘的实现机制,也就是说按照一定的消息量和时间间隔来刷盘,而最终刷新到磁盘的这个动作,是由操作系统来调度的,所以如果在刷盘之前系统崩溃,就会导致数据丢失。
Kafka并没有提供同步刷盘的实现,所以针对这个问题,需要通过Partition的副本机制和acks机制来一起解决。
简单说一下Partition副本机制,它是针对每个数据分区的高可用策略,每个partition副本集包含唯一的一个Leader和多个Follower,Leader专门处理事务类的请求,Follower负责同步Leader的数据。
在这样的一种机制的基础上,kafka提供了一个acks的参数,Producer可以设置acks参数再结合Broker的副本机制来个共同保障数据的可靠性。acks有几个值的选择:
acks=0, 表示producer不需要等Broker的响应,就认为消息发送成功,这种情况会存在消息丢失。
acks=1, 表示Broker中的Leader Partition收到消息以后,不等待其他Follower Partition同步完,就给Producer返回确认,这种情况下Leader Partition挂了,会存在数据丢失。
acks=-1,表示Broker中的Leader Parititon收到消息后,并且等待ISR列表中的follower同步完成,再给Producer返回确认,这个配置可以保证数据的可靠性。
最后,就是Consumer必须要能消费到这个消息,实际上,我认为,只要producer和broker的消息可靠的到了保障,那么消费端是不太可能出现消息无法消费的问题,除非是Consumer没有消费完这个消息就直接提交了,但是即便是这个情况,也可以通过调整offset的值来重新消费
2. Kafka 如何避免重复消费
首先,Kafka Broker上存储的消息,都有一个Offset标记。然后kafka的消费者是通过offSet标记来维护当前已经消费的数据,每消费一批数据,Kafka Broker就会更新OffSet的值,避免重复消费。默认情况下,消息消费完以后,会自动提交Offset的值,避免重复消费。
Kafka消费端的自动提交(enable.auto.commit 默认值:true)逻辑有一个默认的5秒间隔(auto.commit.interval.ms 默认值:5000毫秒),也就是说在5秒之后的下一次向Broker拉取消息的时候提交。
所以在Consumer消费的过程中,应用程序被强制kill掉或者宕机,可能会导致Offset没提交,从而产生重复消费的问题。
在Kafka里面有一个Partition Balance机制,就是把多个Partition均衡的分配给多个消费者。
Consumer端会从分配的Partition里面去消费消息,如果Consumer在默认的5秒内没办法处理完这一批消息,就会触发Kafka的Rebalance机制,从而导致Offset自动提交失败。
而在重新Rebalance之后,Consumer还是会从之前没提交的Offset位置开始消费,也会导致消息重复消费的问题。
基于这样的背景下,我认为解决重复消费消息问题的方法有几个。
- 提高消费端的处理性能,避免触发Balance,比如可以用异步的方式来处理消息,缩短单个消息消费的时间。或者还可以调整消息处理的超时时间。还可以减少一次性从Broker上拉取数据的条数。
- 可以针对消息生成md5,然后保存到mysql或者redis里面,在处理消息之前先去mysql或者redis里面判断是否已经消费过。这个方案其实就是利用幂等性的思想。
3. 什么是 ISR,为什么需要引入 ISR
首先,发送到Kafka Broker上的消息,最终是以Partition的物理形态来存储到磁盘上的。
而Kafka为了保证Parititon的可靠性,提供了Paritition的副本机制,然后在这些Partition副本集里面,存在Leader Partition和Flollower Partition。生产者发送过来的消息,会先存到Leader Partition里面,然后再把消息复制到Follower Partition,这样设计的好处就是一旦Leader Partition所在的节点挂了,可以重新从剩余的Partition副本里面选举出新的Leader。然后消费者可以继续从新的Leader Partition里面获取未消费的数据。
在Partition多副本设计的方案里面,有两个很关键的需求。
• 副本数据的同步
• 新Leader的选举
这两个需求都需要涉及到网络通信,Kafka为了避免网络通信延迟带来的性能问题,以及尽可能的保证新选举出来的Leader Partition里面的数据是最新的,所以设计了ISR这样一个方案。
ISR全称是 in-sync replica,它是一个集合列表,里面保存的是和Leader Parition节点数据最接近的Follower Partition。如果某个Follower Partition里面的数据落后Leader太多,就会被剔除ISR列表。
简单来说,ISR列表里面的节点,同步的数据一定是最新的,所以后续的Leader选举,只需要从ISR列表里面筛选就行了。
所以,我认为引入ISR这个方案的原因有两个:
1. 尽可能的保证数据同步的效率,因为同步效率不高的节点都会被踢出ISR列表。
2. 避免数据的丢失,因为ISR里面的节点数据是和Leader副本最接近的。
4. Kafka 如何保证消息消费的顺序性
首先,在kafka的架构里面,用到了Partition分区机制来实现消息的物理存储,在同一个topic下面,可以维护多个partition来实现消息的分片。
生产者在发送消息的时候,会根据消息的key进行取模,来决定把当前消息存储到哪个partition里面,并且消息是按照先后顺序有序存储到partition里面的。
在这种情况下,假设有一个topic存在三个partition,而消息正好被路由到三个独立的partition里面。然后消费端有三个消费者通过balance机制分别指派了对应消费分区。因为消费者是完全独立的网络节点,所有可能会出现消息的消费顺序不是按照发送顺序来实现的,从而导致乱序的问题。
-
针对这个问题,一般的解决办法就是自定义消息分区路由的算法,然后把指定的key都发送到同一个Partition里面,接着指定一个消费者专门来消费某个分区的数据,这样就能保证消息的顺序消费了。

-
另外,有些设计方案里面,在消费端会采用异步线程的方式来消费数据来提高消息的处理效率,那这种情况下,因为每个线程的消息处理效率是不同的,所以即便是采用单个分区的存储和消费也可能会出现无序问题,
针对这个问题的解决办法就是在消费者这边使用一个阻塞队列,把获取到的消息先保存到阻塞队列里面,然后异步线程从阻塞队列里面去获取消息来消费。
5. 如何处理消息队列的消息积压问题
通常来说,消息积压的原因是生产者的消息生产速度大于消费者的消费速度,遇到这个问题的时候,需要排查具体的原因再提出解决方案。
如果当前不是因为系统bug导致的,那我们可以优化消费端的逻辑,比如通过异步的方式来处理消息、或者通过批量处理的方式来消费。
如果通过这两种优化方式还没有缓解,可以考虑对消费端进行水平扩容,从而扩大消费端的消费能力。
如果是因为系统bug导致大量消息堆积,那么首先需要解决系统bug,然后临时做紧急扩容来完成大量消息的消费。
- 首先解决消费端的bug,来保证消费端的正常消息处理工作。
- 接着把现在所有的消费端停止,然后新建一个Topic,然后把Partition分区数量调整成原来的10倍。
- 接着写一个用来实现数据分发的Consumer程序,这个程序专门去消费现在积压的数据,消费后不做处理,而是直接再把这些数据写入临时建立的Topic的10个Partition中。
- 然后临时增加10倍的消费者节点来部署Consumer,专门来消费临时的Partition分区数据。
通过上面这种方法,可以快速把现在堆积的消息处理完。等积压的消息处理结束后,再把恢复成原来的部署架构,把临时的Topic和临时申请的机器释放掉。
6. Kafka消息队列怎么保证exactlyOnce,怎么实现顺序消费
在回答这个问题之前,先来了解一下Kafka的运行机制:
当我们向某个Topic发送消息的时候,在Kafka的Broker上,会通过Partition分区的机制来实现消息的物理存储。
一个Topic可以有多个Partition,相当于把一个Topic里面的N个消息数据进行分片存储。消费端去消费消息的时候,会从指定的Partition中去获取。
在同一个消费组中,一个消费者可以消费多个Partition中的数据。但是消费者的数量只能小于或者等于Partition分区数量。
理解了Kafka的工作机制以后,再来理解一下exactlyOnce的意思,在MQ的消息投递的语义有三种:
• At Most Once: 消息投递至多一次,可能会丢但不会出现重复。
• At Least Once: 消息投递至少一次,可能会出现重复但不会丢。
• Exactly Once: 消息投递正好一次,不会出现重复也不会丢。
我们只能通过一些其他手段来达到Exactly Once的效果。也就是确保生产者只发送一次,消费端只接受一次。
- 生产者可以采用事务消息的方式,事务可以支持多分区的数据完整性,原子性。并且支持跨会话的exactly once处理语义,即使producer宕机重启,依旧能保证数据只处理一次。
开启事务首先需要开启幂等性,即设置enable.idempotence为true。然后对producer消息发送做事务控制。如果出现导致生产者重试的错误,同样的消息,仍由同样的生产者发送多次,这个消息只被写到 Kafka broker 的日志中一次。 - 虽然生产者能保证在Kafka broker上只记录唯一一条消息,但是由于网络延迟的存在,有可能会导致Broker在投递消息给消费者的时候,触发重试导致投递多次。所以消费端,可以采用幂等性的机制来避免重试带来的重复消费问题。
其次,关于实现顺序消费问题。
在Kafka里面,每个Partition分区的消息本身就是按照顺序存储的。所以只需要针对Topic设置一个Partition,这样就保证了所有消息都写入到这一个Partition中。而消费者这边只需要消费这个分区,就可以实现消息的顺序消费处理。
7. 说一下Kafka中Partition分区副本的Leader选举算法
在Kafka的架构中,一个Topic逻辑主题,可以分成多个Partition分区实现消息内容的物理存储。同时,为了保证Partition分区的可靠性,Kafka设计了分区副本的概念,也就是一个Partition可以设置多个副本。在多个副本中,由于设计到数据的同步,所以Kafka针对Partition分区副本集,设置了Leader副本和Follower副本。Leader副本负责处理所有的读写请求,Follower副本只负责从Leader副本同步数据。
Kafka首先会选择一个具有最新数据的副本作为新的Leader,也就是ISR集合中的副本。
其中,ISR(In-Sync Replica)是指与Leader同步的副本集合,它们的数据同步状态与Leader最接近,并且它们与Leader副本的网络通信延迟最小。如果ISR集合中没有可用的副本,Kafka会从所有副本中选择一个具有最新数据的副本作为新的Leader。
在这种情况下选举出来的Leader,由于和原来老的Leader节点的数据存在较大的延迟,会造成数据丢失的情况。
所以Kafka设计者把这个功能开关的选择交给了开发者,如果愿意接受这种情况,可以通过unclean.leader.election.enable参数来设置。
开启之后虽然会造成数据丢失,但是至少可以保证依然能对外提供服务,保证了可用性。
8. Kafka中一个Topic有三个Partition,同一个消费组中两个消费者如何消费
先来了解一下Kafka的运行机制:
当我们向某个Topic发送消息的时候,在Kafka的Broker上,会通过Partition分区的机制来实现消息的物理存储。
一个Topic可以有多个Partition,相当于把一个Topic里面的N个消息数据进行分片存储。
消费端去消费消息的时候,会从指定的Partition中去获取。
在同一个消费组中,一个消费者可以消费多个Partition中的数据。但是消费者的数量只能小于或者等于Partition分区数量。
而这里提出来的问题是,一个消费组中两个消费者去消费三个Partition,很自然的想到其中一个消费者需要消费两个Partition。
这个问题涉及到Kafka里面的Consumer Group Coordinator ,也就是消费组协调器。
它会根据消费者订阅的Topic中的partition数量、
和消费组中的消费者实例数量来决定每个消费者消费哪些Partition。
这个算法会在消费组中选择一个消费者实例作为Leader,
Leader负责分配Partition给消费者实例,并协调消费者实例之间的Partition分配和Reblance
当一个消费者实例加入或离开消费组的时候,协调器会触发Partition的重新分配,确保所有Partition都能被消费者实例均匀地消费。
Kafka还提供了三种Partition分配策略,
- Round-robin(轮询),它会将Partition均匀地分配给消费者实例。
- Range(范围),它会按照Partition的范围进行分配。
- Sticky(粘滞分配),它会尽可能地将同一Partition分配给同一个消费者实例。
更多推荐




所有评论(0)