跳过正文
  1. 全部/
  2. 笔记/
  3. 面试题/

Kafka

目录

1.应用场景
#

1.1 消息队列常见应用场景有哪些?
#

  • 系统解耦:在重要操作完成后,发送消息到Kafka中,由别的服务系统来消费消息完成其他操作​
  • 流量削峰:一般用于秒杀或抢购活动中,来缓冲系统短时间内高流量带来的压力​
  • 异步处理:通过异步处理机制,可以把一个消息放入队列中,但不立即处理它,在需要的时候再进行处理​
  • 消息分发:比如一条消息需要传递给多个服务,这时候就可以用Kafka进行分发

1.2 什么情况下需要解耦?
#

比如发送短信场景,模块A发送消息给B,模块B发送短信给客户,A不需要得到回应,对于A而言只需要触达B就行了,这时候就可以引入消息队列,A将消息投递到了消息队列,B自己去处理,A不用再关心。这样可以收获更强的稳定性以及接口性能。

1.3 什么情况下需要削峰?
#

比如一个模块B一瞬间收到100个请求,如果B的承压能力非常差,或者B有什么资源限制,那么这100个请求下来,B可能就挂了或者报错了,这时候就可以引入消息队列,前置服务先将消息打到消息队列,B再根据自己的消费能力逐步消化。

1.4 什么情况下需要分发消息?
#

比如信息更新场景,某个用户信息更新了,而B、C、D三个模块都需要缓存这个信息,那么用户信息更新之后就可以发一条信息到消息队列,B、C、D只要订阅了相关主题,就都可以获得这条信息并做对应的缓存更新。

1.5 在项目开发中,你会怎么选择消息队列?
#

我会对比常见的几个消息队列,从功能性、性能、健壮性上去做对比,比如如果我们的性能要求非常高,我就会偏向于Kafka和RocketMQ这类性能又高,又自带扩展性的消息队列。​

当然还有一个考虑点就是团队技术栈,如果基本要求都满足,会更倾向于选择团队使用成熟的组件。

1.6 Kafka和RocketMQ有什么区别
#

我之前主要都是在研究/使用Kafka,RocketMQ只有一些粗浅的认知,比如RocketMQ可以支持更多的主题,适合于需要非常多主题的场景,以及RocketMQ的功能应该更完善一些,比如天然支持死信队列

1.7 Kafka的优势是什么?
#

Kafka优点非常多,我认为最核心的是

  • 高吞吐
  • 高可靠
  • 天然支持分片水平扩展

另外还有一个非常重要的考量就是Kafka多年来被广泛使用,并且社区非常活跃,这就意味着它是经过市场验证的好产品。​

2.服务端
#

2.1 Kafka的大概框架是怎么样的
#

我们可以先把Kafka宏观看作三层:Producer,Server,Consumer,生产者发送消息,服务端负责存储消息,消费者负责拉取消息。 其中服务端其实就是由多个Broker节点组成,而我们平常说的主题就是在Broker节点上 Topic是个逻辑概念,实际物理存储形式是主题分片,也就是Partition

2.2 如何获取topic主题的列表
#

  • Kafka是提供了获取主题列表的接口的,可以使用kafka-topics.sh 这个工具获取
  • 如果要在我们业务服务来获取的话,那主流语言都是支持获取的,比如Java和Golang,都可以用KafkaAdminClient直接调用接口来获取。

2.3 有了Topic,为何还要Partition (Kafka 为什么要把消息分区)
#

一般而言一个业务可以用一个主题,但是就算分了Topic,单个业务的信息还是可能会非常多,所以需要能进一步进行分治,也就是一个主题由多个Partition组成,这样相同的主题也具备了更高的并发度

2.4 Partition是逻辑概念还是物理概念?
#

Partition是物理概念,因为数据实实在在是写入到Partition文件里去的。而因为设计了Partition的存在,Topic就只是一个逻辑概念了。

2.5 介绍一下分区的分配策略
#

  • Range Assignor:基于范围的分配策略,将分区按照范围分配给消费者。​
  • RoundRobin Assignor:基于轮询的分配策略,分区均匀地分配给消费者。​
  • Sticky Assignor:优先保持当前的分配状态,并尽量减少在再平衡过程中的分区移动。​
  • CooperativeStickyAssignor:和Sticky Assignor的策略是基本一样的,区别在于该协议将原来的一次大的全部分区重平衡,改成多次小规模分区重平衡。简单理解就是渐进式重平衡。

2.6 Kafka 创建 Topic 时如何将分区放置到不同的 Broker 中
#

一个Partition只对应一个Broker,一个Broker可以存放多个Partition。 kafka是先随机挑选一个broker放置分区0,然后再按顺序放置其他分区

2.7 消息存入哪个Partition的规则 / Kafka 中分区分配的原则
#

  1. 如果指定了Partition,那么就是发送到特定的Partition​
  2. 如果没有指定Partition,但是指定了一个Key,那么就是根据Key的Hash对Partition数目取模来决定是哪个Partition,也就是说只要发送时指定了相同的Key,那么相关消息一定会发送到相同的Partition。​
  3. 如果没有指定Partition,也没有指定Key,那么就采取轮询调度算法,也就是每一次把来自用户的请求轮流分配给Partition。

2.8 Kafka服务端可接收的消息最大默认多少字节,如何修改
#

Kafka可以接收的最大消息默认为1MB

2.9 Kafka 的Topic中 Partition 数据是怎么存储到磁盘的
#

Topic 中的多个 Partition 以文件夹的形式保存到 Broker,每个分区序号从0递增,且消息有序。 Partition 文件下有多个Segment(xxx.index,xxx.log),Segment文件的大小是可以进行配置的。默认为1GB。如果大小大于1GB时,会滚动一个新的Segment并且以上一个Segment最后一条消息的偏移量命名。

2.10 Kafka如何清理数据/Kafka数据越积越多怎么办
#

可以用基于时间的保留策略,这种策略允许用户指定消息的保留时间(如 7 天)。超过指定时间的消息将被自动删除。​

也可以用基于大小的保留策略,Kafka 允许用户指定日志的最大尺寸。一旦日志的大小超过了配置的值,Kafka 将开始删除最早的消息。

3.生产者
#

3.1 介绍一下生产消息的流程
#

第一步是构建消息,即将要发送的内容,打包成一个Kafka的消息结构;​ 第二步是序列化消息为二进制内容,以在网络中传输​ 第三步是进行分区选择,即计算要发到哪个Partition,发送消息到该Partition对应的Broker

3.2 讲一讲kafka的ack的三种机制
#

第一种模式,ack=0,生产者在发送消息后不会等待来自服务器的确认,所以生产者实际是不知道消息是否成功,也就无从去重试,生产可靠性是最低的。​

第二种模式,ack=1,生产者会在消息发送后等待主节点的确认,但不会等待所有副本的确认。​

第三种模式,ack=all,只有在所有副本都成功写入消息后,生产者才会收到确认。这确保了消息的可靠性,但延迟显然是最高的

3.3 生产过程中何时会发生QueueFullExpection以及如何处理
#

等待重试、增加Kafka的缓冲区大小、限流控制

3.4 Kafka 生产者何时发出消息
#

  • 累计的数据大小达到Batch大小,默认16KB​
  • 缓冲区中累计的空闲等待时间间隔,默认0ms,也就是默认收到数据就发送​

Kafka 生产者调用 send() 后并不会立即发送消息,而是先写入本地内存缓冲区 RecordAccumulator,由后台 Sender 线程根据 batch.size、linger.ms 或 buffer.memory 等条件批量发送到 Broker,从而实现高吞吐的异步批处理模型。

3.5 生产者发送消息的模式有几种
#

Kafka Go 生产者主要有三种写法: 同步发送(SendMessage 阻塞等待 ACK) 异步发送(通过 channel 或 callback 写入内部队列由后台 goroutine 批量发送) 以及事务型发送(支持 exactly once 语义)

4.消费者
#

4.1 Kafka 消费者是推模式还是拉模式/Kafka 消息的消费模式
#

Kafka采用拉模式拉取消息,采用拉模式可以使每个消费者以自身的消费能力去消费。 拉模式有个缺点是,如果Broker没有可供消费的消息,将导致Consumer不断在循环中轮询,直到新消息到达。 为了避免这点,kafka消费者可以使用在消费数据时传入timeout参数,在这个时间范围内进行阻塞等待,直到有数据或超时后再返回。

4.2 消费者故障,出现活锁问题如何解决?
#

可以使用最大拉取间隔这个参数来解决活锁问题,即max.poll.interval.ms,如果消费者轮询间隔大于了这个值,消费者就会离开分区,这样其它消费者就可以接管对应分区。

4.3 有消费者为什么还要消费者组?
#

消费者组:一组共同消费同一个 Topic 的消费者实例集合,由 Kafka 保证同一 partition 在同一时刻只会被组内一个消费者消费。

  1. 对主题分片的分配问题,让每个分片都能有消费者处理,又不至于所有消费者处理同一个分片​
  2. 面对主题分片的变化,消费组可以自动调整,也就是再平衡​
  3. 对于业务开发者而言,有了消费组,就只用关心主题维度,而不用关心分片维度,很大程度降低了理解和应用难度

4.4 介绍一下再平衡机制
#

消费者组再平衡是一个关键机制,用于管理和分配主题分区给消费者组中的各个消费者。再平衡过程可以确保数据负载在消费者之间均匀分布,并在消费者加入或离开时自动调整分区的分配。我记得是有几个分区策略可以选择的,​

分别是范围分配,轮询分配,粘性分配,合作粘性分配,其中合作粘性分配和粘性分配一样都是尽可能减少变动,不同点是合作粘性分配下,是把大的分区平衡分为多次小规模的分区平衡,以尽可能减少影响。

4.5 Kafka什么情况下会Rebalance
#

新消费者加入:当一个新的消费者加入消费者组时,Kafka需要重新分配分区,以包括新的消费者。​

消费者离开:当一个消费者离开(无论是正常关闭还是崩溃)时,需要重新分配该消费者负责的分区给其他消费者。​

主题分区变化:当主题的分区数量发生变化时(例如,增加新的分区),Kafka需要重新分配这些分区。

4.6 Rebalance有什么影响
#

  1. 重复消费,如果某个消费者离开消费组时还没来得及提交Offset,当再平衡之后,接盘对应分区的消费者就会重复消费,浪费资源。​

  2. 性能变差,上面介绍了再平衡是需要相对复杂的流程去实施的,在实施再平衡的这个过程中,消费速度也会受到影响

4.7 介绍一下重平衡的具体执行流程
#

首先是暂停消费,其作用是防止在重新分配期间发生数据丢失或重复,接着由消费组协调器触发再平衡,进而重新分配分区,最后就是开启消费,简单来说就是通知消费者,然后消费者就可以恢复消费了。

4.8 介绍一下重平衡时的分区策略
#

第一类是急切再平衡(Eager Rebalance)。Range Assignor、RoundRobin Assignor、Sticky Assignor都属于Eager Rebalance,可以理解为急切的再平衡,因为它太过急切,所以顾不得那么多,为了快点搞定这件事,一旦开启再平衡所有消费者都会停止从 Kafka 消费并放弃其分区的成员资格。​

第二类是增量再平衡(Incremental Rebalance),CooperativeStickyAssignor这个2.3版本之后引入的优化策略就属于这一类,在此模式下,只有部分分区会从某个消费者移动到另外一个消费者,其它不受重新平衡影响的 Kafka 消费者可以继续处理数据而不会中断。

4.9 解释下Kafka中位移(offset)的作用
#

每条消息在Kafka中会有Partition ID以及OFFSET,通过这两个信息就可以定位到一条消息。消费者组消费消息之后会提交它在某个Partition对应的OFFSET,这样子下一次就可以从这个位置开始消费。

4.10 如何控制消费的位置
#

每条消息在Kafka中会有Partition ID以及OFFSET,通过这两个信息就可以定位到一条消息。消费者组消费消息之后会提交它在某个Partition对应的OFFSET,这样子下一次就可以从这个位置开始消费。

4.11 Consumer默认从哪里开始消费
#

如果没有提交过偏移,那么会根据auto.offset.reset这个配置决定从哪里开始消费,默认是latest,消费者将从当前最新的数据开始读取,如果业务需要,还可以选择从最早的偏移开始读取。

4.12 Consumer怎么手动指定开始消费偏移
#

Golang客户端有方法能指定开始消费的Offset,我之前有使用过confluent的kafka库,方法名应该就是Assign,将希望的偏移量传入这个方法即可。

4.13Kafka消费消息是推还是拉?
#

Kafka的消费者使用的拉模式来获取信息,也就是说每次消费者是发消息到Kafka的Broker来获取信息,而不是由Kafka的Broker主动推送。​

选择拉模式的主要原因还是为了让消费者可以按自身情况来控制消费速度,根据系统资源利用情况(如 CPU、内存等)、业务需要等因素合理拉取消息,避免因消息处理速度不合理带来的资源浪费或过载。

4.14 Kafka消费者提交之后就会清理掉数据吗
#

在Kafka中,如果消息被消费者消费并提交了对应偏移,这条消息不会被删除,可以通过更改该消费者的偏移再次消费,也可以被其它消费者消费

5.实践经验
#

5.1 Kafka 中什么情况下会出现消息丢失的问题
#

  • 消息生产时如果设置的acks不是全部副本,那么如果在follower副本未完成同步之前,leader副本挂掉了,消息就会丢失。
  • 存储时如果没有用多副本备份,消息也可能会丢失。
  • 最后就是消费时,如果没有确认消费成功再提交offset,而这时候消费者又挂掉了,那么消息同样会丢失。

5.2 Kafka 如何保证消息不丢失
#

从消息流转环节分析,分别考虑生产环节、存储环节、消费环节来看。

  • 首先为主题分区配置好多副本,
  • 并且设置写入acks参数为全部副本,
  • 最后就是消费时候一定要确认消费成功再提交offset,这样即使消费者挂掉了,重启之后也能拉到原来那条未成功消费的消息。

5.3 MQ消息积压了怎么办
#

扩容、降级、排查异常

  • 如果分区数大于消费者数量,那么通过扩容消费端的实例数来提升总体的消费能力;如果相等,那么既需要扩容消费者数量同时需要扩容分区数。​
  • 如果短时间内没有足够的服务器资源进行扩容,可以考虑将系统降级,通过关闭一些不重要的分支业务,让系统还能正常运转,服务一些重要业务。​
  • 还有一种不太常见的情况,你通过监控发现,无论是发送消息的速度还是消费消息的速度和原来都没什么变化,这时候你需要检查一下你的消费端,是不是消费失败导致的一条消息反复消费这种情况比较多

5.4 Kafka 如何保证消息不重复消费
#

消费逻辑需要是幂等的,保证不产生重复影响,实现方式很多,比如MySQL设置唯一索引、额外使用记录表来判重等方式。

5.5 Kafka 如何实现精准一次语义
#

本质就是不重复+不丢失。 不重复的核心是幂等消费。 不丢失的核心是为主题分区配置好多副本,并且设置写入acks参数为全部副本,同时消费时候一定要确认消费成功再提交offset,这样即使消费者挂掉了,重启之后也能拉到原来那条未成功消费的消息。

5.6 假设你有个业务希望进入Kafka的消息都是有序的,你会怎么做?
#

如果指定了Partition,那么就是发送到特定的Partition;如果没有指定Partition,但是指定了一个Key,那么就是根据Key的Hash取模来决定是哪个Partition;如果都没有指定,就是依次轮替着写入。​

所以我们可以用一个能标识业务的唯一名字来当Key,比如秒杀,就叫SecKill,指定Key之后算出来一定是落在相同的Partition,也就保证了顺序。

6.高可用(副本)
#

6.1 Replica、Leader 和 Follower 三者的概念
#

Replica:Replica是指Kafka集群中的一个副本,它可以是Leader副本或者Follower副本的一种。每个分区都有多个副本,其中一个是Leader副本,其余的是Follower副本。每个副本都保存了分区的完整数据,以保证数据的可靠性和高可用性。​

Leader:Leader是指Kafka集群中的一个分区副本,它负责处理该分区的所有读写请求。Leader副本是唯一可以向分区写入数据的副本,它将写入的数据同步到所有的Follower副本中,以保证数据的可靠性和一致性。​

Follower:Follower是指Kafka集群中的一个分区副本,Follower副本不能直接向分区写入数据,它只能从Leader副本中复制数据,并将数据同步到本地的副本中,以保证数据的可靠性和一致性。在leader副本挂掉的时候,follower副本有机会被选举为新的leader副本从而保证分区的可用性。

6.2 Kafka 中 AR、ISR、OSR 三者的概念
#

AR(Assigned Replicas):AR是指分区的所有副本,包括Leader副本和Follower副本。​

ISR(In-Sync Replicas):ISR是指与Leader副本保持同步的副本集合。ISR中的副本与Leader副本保持同步,即它们已经复制了Leader副本中的所有数据,并且与Leader副本之间的数据差异不超过一定的阈值(Follower副本能够落后Leader副本的最长时间间隔)。并且ISR副本集合是动态变化的,不是一成不变的。ISR中的副本可以被选举为新的Leader副本,以保证分区的正常运行。​

OSR(Out-of-Sync Replicas):OSR是指与Leader副本不同步的副本集合。OSR中的副本与Leader副本之间的数据差异超过了一定的阈值,或者它们还没有复制Leader副本中的所有数据。OSR中的副本不能被选举为新的Leader副本,除非开启了Unclean选举。它们只能等待与Leader副本同步,或者被替换为新的副本。​

6.3 分区副本什么情况下会从 ISR 中剔出
#

每个Partition都会由Leader 动态维护一个与自己基本保持同步的ISR列表。所谓动态维护,就是说如果一个Follower比一个Leader落后超过了给定阈值,默认是10s,则Leader将其从ISR中移除。如果OSR列表内的Follower副本与Leader副本保持了同步,那么就将其添加到ISR列表当中。

6.4 分区副本中的 Leader 如果宕机但 ISR 却为空该如何处理
#

可以通过unclean选举配置参数来决定是否从OSR中选举出leader:​ 如果是该参数是true:允许 OSR 成为 Leader,但是 OSR 的消息较为滞后,可能会出现消息丢失的问题​ 否则:坚决不能让那些OSR竞选Leader。这样做的后果是这个分区就不可用了。

6.5 分区副本之间同步,是推还是拉 (kafka如何主从同步)
#

数据是先写入到Leader副本,同步时候是Follower副本去主动拉取消息,拉的优势在于副本机器可以根据自身的负载情况来拉取。

6.6 高可用机制是怎么实现的
#

Kafka天然支持多副本机制,每个副本都有完整的数据,这些副本分散在不同的Broker上,就算主副本所在Broker的磁盘损坏了,其它Broker也能把数据找回来并升级为主副本,通过这种方式Kafka就实现了高可用。

6.7 Kafka是怎么为分片选择主副本的
#

Kafka维护了一个叫ISR的列表,ISR里的副本都是包含完整数据的,当没有Leader,或原有Leader挂掉了,Kafka就会从ISR列表中选择第一个副本升级为Leader。

6.8 Kafka怎么知道Leader挂了
#

Broker中会选出一个来担任Controller,选Controller负责监测Leader的状态,这样Leader挂掉之后Controller就能感知到,并从ISR选择出新的Leader。

7.高性能
#

7.1 Kafka为什么这么快/Kafka性能为什么这么高/Kafka吞吐量为什么这么大
#

我知道有几个点都有提升Kafka处理的速度,包括顺序写、零拷贝、数据压缩、批量操作(说自己比较熟悉的,多一个少一个,无所谓),其中我觉得影响最大或者说最具Kafka特色的,就是顺序写,通过磁盘的顺序写入,非常直接优化了兼顾了性能和复杂度,其次我印象最深的是批量操作和数据压缩,这两个也是业务优化的常见思路。​

7.2 聊聊Kafka顺序写机制
#

顺序写指的是按顺序将数据写入磁盘,对于Kafka而言,其实就是直接追加到磁盘文件末尾,顺序写性能非常高,甚至接近内存写,所以出于Kafka自身定位、复杂度、性能等综合考虑,Kafka选择了顺序写。

7.3 聊聊Kafka页缓存机制
#

Page Cache可以简单看作热点磁盘数据的内存缓存,当消息写入时,是先写入Page Cache,后面由操作系统将其刷入磁盘,这样性能就会提升很多,同时,如果查询时候发现PageCache中有对应数据,那么也就不用去磁盘读取。​

值得一提的是,Kafka是生产消费者模式,即生产了消息,在无积压情况下,这个消息很快就会被消费,也就是说我们生产消费时写入了Page Cache,而很快就有消费者来触发Kafka应用程序读取对应数据,而这个时间间隔很短,PageCache命中的可能性会很高,自然提升效果就会非常大了。

7.4 聊聊Kafka零拷贝机制
#

零拷贝是一种 I/O 优化技术,其核心是减少用户态与内核态之间的数据拷贝次数,通过 mmap、sendfile 等机制让数据尽可能在内核空间完成传输,从而降低 CPU 开销并提升系统吞吐量,Kafka 和 Nginx 等高性能系统广泛使用该技术。

7.5 聊聊Kafka分层设计机制
#

分层设计其实就是出自分治思考,对于Kafka而言,所有消息写入一个文件,那肯定扛不住,所以提出了Topic的概念,可以一个业务写入一个Topic,单个Topic不具备扩展性,扛不住大流量业务,所以Topic又进行了分片,也就是Partition,一个业务的消息可以根据规则写入多个Partition,单个Partition是不是也需要能扩展,于是一个Partition又可以切分为多个segment文件,segment文件支持按需滚动增长,所以Partition就具备了扩展性。这就是Kafka一以贯之的分层设计。

7.6 Kafka 文件高效存储设计原理
#

Kafka把topic中一个大的parition文件分成多个小的segment file,通过多个segment file,就容易定期清除或删除已经消费完的文件,减少磁盘的占用。​

  • 通过索引元数据来管理消息的位置和偏移量,以便快速定位和读取消息。​
  • 通过索引元数据全部映射到内存,可以避免segment file的IO磁盘操作。​
  • 通过索引文件稀疏存储,可以大幅降低索引文件元数据占用空间大小。

7.7 聊聊Kafka哪些环节用了批量操作
#

Kafka主要有2个批量操作的地方,一个是批量生产,也就是批量发送,其实就是通过发送缓冲,将数据缓冲起来,等聚集了一批数据,再一次性发送给Broker。另一个是批量消费,本质就是一次多拉几条消息,一起消费。Kafka服务层也会将多条消息一次性批量写进磁盘里以提高性能。

7.8 聊聊Kafka数据压缩
#

通过压缩可以让传输的数据变小,以节约带宽,所以我们可以通过压缩消息来提升Kafka性能,一般而言就是在发送方进行压缩,有时候也可以在Broker侧进行压缩。压缩适用于CPU比较富裕,带宽相对不足的情况,而消息的传输大多数是符合这个情况的。

8.扩展
#

8.1 Zookeeper对于Kafka的作用是什么
#

Zookeeper拥有分布式协调能力,Kafka主要是用Zookeeper来管理Broker/Topic数据、存储配置、选择Controller,曾经也会存储消费者偏移信息

8.2 Kafka你用的是哪个版本?
#

在xxx项目中用的是Kafka 2.5.0版本,主要原因是这个版本团队用了很多年都是非常稳定的,所以新老项目都一直延用这个版本。

Reply by Email