Java面试题-kafka
1.讲一下kafka
kafka是一个分布式流处理平台,也是一个消息中间件。kafka基于生产者-消费者模型,可以实现应用的解耦。主要有以下几个核心概念
- 消息(message):kafka的数据单位。
- 主题(topic):kafka中的消息都是以主题进行归类的。
- 分区(partition):一个主题可以包含多个分区,每一个分区都是有序且不可变的消息序列。
- 生产者(producer):生产者向kafka发送消息,负责将数据写入kafka。
- 消费者(consumer):消费者从kafka获取消息,负责从kafka读取消息。
- broker:一个kafka节点就是一个broker。
- offset:消息在分区中的位置索引。
2.kafka为什么需要分区?
- 负载均衡:消息通过指定的分配策略分配到不同的broker上,避免单个broker成为性能瓶颈。
- 容错性:kafka会将分区数据复制到多个broker上,提高系统容错性,保证数据的可用性。
- 顺序消费:生产者按照顺序发送到同一个分区的消息,消费者也会按照这个顺序进行消费。
- 可扩展性:当主题数量不断增长时,可以增加分区分散负载,不需要对系统进行大规模重构。
3.kafka如何保证顺序消费?
- 发送到单个分区:生产者按照顺序将消息发送到单个分区,一个分区的数据只能被一个消费者实例消费。
- 单线程消费:消费者组中只有一个消费者实例,但是实际项目中基本不存在这种情况。
- kafka事务:可以使用kafka的事务,严格的进行生产和消费。但是这会牺牲性能和吞吐量。
4.如何保证kafka的消息发送到指定的分区?
- 在生产消息时直接指定分区号。
- 自定义分区,通过实现 Partitioner 类,自定义逻辑分区。
- 默认的分区策略,基于key进行哈希计算,然后与分区数取模,得到分区号。
5.kafka中,消费者提交消费位移时提交的是当前消费到的最新消息的offset还是offset+1?
在kafka中,消费者提交的是最新的一条消息唯一,也就是offset+1。主要是为了保证消费者在下次启动时,从未消费的消息开始读取,防止重复消息。
6.kafka中,有哪些情况会造成重复消费?
- 没有正确的提交位移会造成重复消费。
- 当消费者组内成员数量发生变化时,会触发再平衡。再平衡的过程中,如果消息已经被处理但是还没有提交位移,就会造成重复消费。
- 生产者在发送消息时,比如网络问题没有收到broker的确认,会重试发送消息。那么消费者在消费时可能就会重复消费。
7.那么如何防止重复消费?
- 手动提交位移:确保消息处理成功后再提交。
- 变更消费者组成员数量时,确保没有消息在消费。
- 配置 enable.idempotence=true ,生产者会自动进行去重,确保消息只被写入一次。
- 通过kafka事务保证原子性,也可以避免消息重复消费。
8.说一下acks参数对消息持久化的影响
- acks = 0:生产者发送消息后不需要等待broker的确认就认为消息已经发送成功。这种情况下如果broker宕机,可能会造成消息的丢失。
- acks = 1:生产者发送消息后会等待broker的确认,broker返回后,消息发送成功。但是broker的返回只是leader副本确认后就返回,所以如果leader副本在同步follower副本时发生故障,也会造成消息丢失。
- acks = -1/all:生产者的消息在所有副本都接收到,才会收到broker的确认。这种设置提供了最高的持久化保证,即使某个broker发生了故障,消息也不会丢失。
9.kafka中,哪些情况会造成消息漏消费?
如果消费者进行消费时发生了异常,没有正确的处理位移,可能会导致消息漏消费。消费者消费能力不足时,会导致消息堆积。当堆积的消息超过一定数量或者时间限制时,会触发清理机制,导致消息被清理继而漏消费。
10.kafka的文件存储机制是怎么样的?
kafka的每个分区对应一个文件夹,文件夹中包含两部分文件:数据文件和索引文件。数据文件保存着实际的消息数据,索引文件保存着消息的偏移量和物理位置的对应关系,可以快速的查找某个消息的位置。
11. kafka如何保证消息的可靠性?
- 分区复制:kafka将每个分区的数据副本分布在多个broker上,即使某个broker宕机,数据仍然能够通过其他broker获取。
- ISR机制:kafka使用ISR机制,只有在所有副本已经同步到最新数据时,才会将消息标记为已提交,这样能够保证所有副本的数据都是一致的。
- 持久化:kafka将消息持久化到磁盘中,即使出现故障,也能够快速恢复。
- 消息重试:kafka允许生产者在发送消息时进行重试,保证消息不会丢失。
- 消费者位移:消费者位移用来记录已经消费的消息位置,确保消息不会被重复消费。
12.kafka的消息传递机制是pull还是push?
kafka采用的是pull模式,也就是消费者从broker中主动拉取数据。这样消费者可以自己控制消费的速度和位置,避免了消息积压和消费者压力过大。并且可以通过设置重复消费一些消息,也保证了消息的一致性和可靠性。
13.kafka如何防止,如果broker没有可供消费的消息,将导致consumer不断在循环中轮询?
消费者在调用poll()方法消费数据时,可以加上timeout参数,当返回空数据的时候,会在Long Polling中进行阻塞,等待timeout再去消费,直到数据到达。
14. kafka中,多个消费者组消费同一个主题,会造成重复消费吗?
- 正常情况下不会,因为每个分区同一时间只会被一个消费者消费。并且因为消费位移的存在,处理完消息之后会提交位移,其他消费者会根据位移位置开始消费,不会重复消费。
- 异常情况下会,比如消费者组重新平衡,kafka会重新分配分区,一些消息可能会被分配给新的消费者,造成重复消费。或者消费时发生异常没有正确的提交位移,也会导致重复消费。
15.kafka为什么性能好、吞吐量大?
- 集群架构:kafka是一个分布式的集群系统,可以将数据分散到不同的节点上进行存储和处理,从而实现横向扩展,提高系统的处理能力。
- 磁盘存储:kafka使用磁盘存储消息,存储数据的容量不再受限于内存的大小。
- 批量发送:kafka可以将多个消息批量发送到broker上,这样可以减少网络传输开销,提高系统吞吐量。
- 零拷贝技术:kafka适合用零拷贝技术来避免数据拷贝的过程,减少了cpu的开销,提供系统性能。
- 压缩算法:kafka支持多种压缩算法,可以对消息进行压缩,减少网络传输开销,提高系统吞吐量。
16.kafka的高可用是怎么实现的?
- kafka分区可以有多个副本,副本分为leader副本和follower副本。
- leader副本负责处理分区的所有读写请求。生产者向分区发送消息时,实际上是发送到leader副本;消费者从分区读取消息时,也是从leader副本读取。
- follower副本则被动地从leader副本复制数据,不处理来自客户端的请求。它们的存在是为了在leader副本出现故障时,能够有一个副本可以被提升为新的leader,继续提供服务。
- follower副本会不断地从leader副本拉取消息,以保持与leader副本的数据同步。这个过程是通过网络通信实现的。
- 当follower副本成功地复制了一条消息后,它会向leader副本发送一个确认消息,表示已经成功复制了该消息。leader副本会跟踪所有副本的同步状态,以确定哪些副本已经跟上了最新的消息。
17.kafka中 leo(log end offset)和hw(high watermark)是什么?
- LEO表示分区中最后一条消息的偏移量,即消息在分区中的存储位置。当 Producer 向 Partition 中写入消息时,LEO 会不断增加,表示消息的写入位置;当 Consumer 从 Partition 中拉取消息时,LEO 会不断变化,表示消息的读取位置。
- HW表示分区中已经被 Consumer 消费的消息偏移量,即消息在 Partition 中的消费位置。当 Consumer 从 Partition 中拉取消息时,HW 会不断增加,表示已经消费的消息的偏移量;当 Consumer 向 Kafka Broker 提交位移时,HW 会被更新为提交的位移,表示消费者已经消费了该位移之前的所有消息。
LEO 和 HW 的关系如下:
1.LEO >= HW:表示 Partition 中还有未被消费的消息,Consumer 可以继续消费这些消息;
2.LEO < HW:表示 Partition 中已经没有未被消费的消息,Consumer 无法继续消费消息,除非有新的消息写入到 Partition 中。
LEO 和 HW 的作用如下:
1.LEO 表示消息的存储位置,可以用于监控 Partition 中消息的写入情况;
2.HW 表示消息的消费位置,可以用于监控 Consumer 消费消息的情况,以及实现消息传输语义的控制。例如,At Least Once 语义中,Consumer 可以将 HW 作为提交位移的参考,避免消息的重复消费。
18.消息被消费后,什么时候会进行删除?没有及时删除会造成空间占满吗?
消息被消费后,不会立马进行删除,而是根据配置的策略进行清除。kafka消息的清除有两种策略,一种基于时间,一种基于分区大小。没有及时删除,当生产者的速度远大于消费者速度时,可能会造成空间占满。
19.kafka默认的主题有哪些?分别做什么用的?
- __consumer_offsets:存储消费组的偏移量信息,每个消费组消费到每个分区的位置。
- __transaction_state:支持 Kafka 事务功能,存储事务的状态信息,比如如事务开始、提交、中止等。
20.kafka为什么要抛弃zookeeper?
kafka在3.3.0版本(2022 年 6 月发布)中正式抛弃了zookeeper,使用KRaft模式进行替代。
zookeeper本质时一个分布式的协调工具,它有一个特点就是,强一致性。当集群中的某个节点发生变化时,需要通知其他节点进行更新,需要等待大多数节点更新完成后才算成功,这样性能就有了瓶颈。当kafka集群很大,分区很多的时候,zookeeper的元数据也会很多,性能就差了。并且zookeeper也是需要选举的,而且发生选举的这段时间不提供服务。
21.kafka消费的时候,消费失败怎么办?
- 延迟重试机制:对于消费失败的消息进行异常捕获,利用指数退避策略进行延迟重试。
- 死信队列:超过最大重试次数后仍然失败,放入死信队列(DLQ)中。
指数退避策略:通过指数级增加重试间隔时间,核心思想是“失败次数越多,下次重试等待时间越长”
- 第一次失败后:1秒
- 第二次失败后:2秒
- 第三次失败后:4秒
- 第四次失败后:8秒
- 第N次失败后:2^n-1秒
22.如何处理死信队列?
- 利用监控工具对死信队列进行监控,及时进行告警。
- 判断消费失败的原因,一般可能有:JSON解析失败、字段类型错误、编码问题等。
- 对消息进行处理后,重新进行消费,并及时的清理死信队列。
23.什么是ISR机制?
ISR机制是kfaka用于保证数据一致性和可靠性的一种机制。
基本原理是:每一个分区都有一个leader副本处理读写,其他的副本都是follower副本,负责数据的备份。follower副本会定期从leader副本拉取最新数据,ISR机制保证只有所有的副本都成功拉取到最新数据后,消息才会被认为“已提交”。对于消费者,因为消费者只能读取已提交的数据,所以ISR机制能确保数据的一致性。
24.什么情况下副本会被踢出ISR?
当follower副本无法在规定时间内拉取到leader副本的最新数据,就会被踢出ISR。
该时间由参数replica.lag.time.max.ms决定,默认30秒。