kafka的架构及常见面试题

发布时间:2026/7/30 23:22:30
kafka的架构及常见面试题 kafka的架构及常见面试题一、介绍Kafka是一种高吞吐量、持久性、分布式的发布订阅的消息队列系统。它最初由LinkedIn(领英)公司发布使用Scala语言编写与2010年12月份开源成为Apache的顶级子项目 。Kafka是一个多分区、多副本且基于zookeeper协调的分布式消息系统。也是一个分布式流式处理平台它以高吞吐、可持久化、可水平扩展、支持流数据处理等多种特性而被广泛使用 。二、架构1生产、消费首先得了解这个比较简单的一个集群图生产者Producer生产消息发送消息的服务消费者Comsumer消费消息处理消息的服务2每一个kafka实例中有什么如上图只画了其中一个具体看看里面是什么broker一个kafka进程就是一个broker也就可以这样理解集群中每一台kafka服务就是broker主题topic在发布订阅的模式下我们需要对消息进行一个区分同一个功能的消息我们发往同一个主题下分区Partition可以看到每一个主题topic下有多个分区。消息会推送到这些分区里面可以增加生产者消费者的对消息生产处理的吞吐量在没有画出来了另外两个kafka中我们会推选出领导Leader以及追随者Follower这个我们后面再聊3简单的消费如上图对于消费者如何消费分片中的消息的其中有下面几点的解释一个Partition只能由一个Consumer来消费一个Consumer可以消费多个不同的Partition。所以我们应该保证每个主题的Partition的数量大于Consumer的数量Consumer越多则吞吐量越高消费得越快。当然要结合第一点Consumer增加或减少时Partition和Consumer的消费关系会自动调整4带group的消费在上一节看到了简单的消费那只不过是同一个group下接下来引入group这个概念上面第三节说的都不错在这里就是要加一个前提条件一个Partition的消息同一个group中的一个Consumer来消费交给了同组的某个Consumer就不能交给同组的其他Consumer了每一个group都可以完整消费主题中的所有消息5消费partition里面的消息如上图有以下几点特性一个Partition内部的消息是有序的越新的消息offset越大。不同partition的消息根据offset无法比较新l旧Consumer顺序地消费partition里的每一条消息可以每读一条就向kafka上报(commit)当前读到了哪个位置(offset)也可以间隔性上报可以读几条再上报offset比如说每读5条上报更新一下offset也可以时间间隔的方式上报offset比如说每隔5s上报更新Consumer重启时kafka根据该group上一次提交的最大offset来决定从哪个地方开始消费。这里就会出现重复消费的问题而解决这个重复消费的问题是面试中的高频问题。不同group之间记录的offset是不同的这也是上一节每个group独立消费topic的消息的原因6生产消息写入Partition关于生产者生产消息至Partition有三种情况按照优先级这样排序生产者可以指定Partition进行写入通过消息携带的key再通过hash分发器计算得到结果来决定去哪个Partition按照时间片轮动选择Partition。比如说当前5分钟往Partition 0中写入下一个5分钟往Partition 1中写入7生产消息写入Partition应答ack在上面一节我们确定了partition存储是哪个接下来还有一个问题就是如果是kafka集群架构的话我们会出现同个Partition有一个Leader多个Follower。在上面确定partition后我们要去寻找它的LeaderLeader partition将消息写入本地磁盘当写入完成后向Producer进行应答响应Follower partition会将消息从Leader那拉回来写入自己的本地磁盘当写入完成后向Leader进行应答响应当leader收到所有的Follower应答后再向Producer应答那么在此刻生产消息的应答ack有三种策略完全不管ack应答Producer只需要Leader Partition应答即可不用管Follower Partition是否写入成功partition需要保证所有的Follower才进行应答8Partition备份机制在kafka集群中我们有Partition的备份机制如下同一个主题下集群中的每个broker都会维护自己的Partition。其中他们会选出Leader、Follower生产者的数据优先推送给Leader每一个Partition都有自己的Leader同一个Topic下的不同Partition尽量分布在不同的broker当有leader的broker宕机后kafka集群会重新竞选那台broker上原本是leader的Partition和下面ISR队列有关。9消息的磁盘存储文件结构分区Partition一个Topic中有多个Partition可以有效地避免了消息的堆积分段segment消息在Partition里面消息是分段来进行存储的每次操作的消息读写都是针对segmengt一个segment包括一个log文件两个index文件三个文件成套出现。前面数字的文件名代表着offset偏移量开始索引位置000000101.log存储具体消息的数据文件000000101.index存储Consumer的offset便宜量的索引文件000000101.timeindex存储消息时间戳的索引文件索引indexkafka分段后的数据建立的索引文件如下图可以看到上面有两个索引文件index文件是记录offset消息和log文件中消息位置的映射关系的文件timeindex文件时记录时间戳和offset关系的文件请注意这边的索引并不会记录每一条消息的索引而是采取稀疏索引也就是隔一段消息才会记录消息的索引。这个消息索引的稠密程度影响kafka存储读取的速度索引越稠密则读取的速度越快索引越稀疏则文件存储的空间越大由于上面存储文件都是采用offset偏移量来命名所以kafka会采取二分查找方法可以大大提交检索效率。三、面试题1如何避免kafka消息丢失1.1出现消息丢失的原因从上面架构上来看kafka丢失消息的原因主要可以分为下面几个场景Producer在把消息发送给kafka集群时中间网络出现问题导致消息无法到达网络抖动原因Producer消息超出大小限制broker收到以后没法进行存储kafka集群接收到消息后保存消息至本地磁盘出现异常集群接收到数据后会将数据进行持久化存储到磁盘消息都是先写入到页缓存然后由操作系统负责具体的刷盘任务或者使用fsync强制刷盘。如果此时Broker宕机且选举了一个落后leader副本限多的follower副本成为新的leader副本那么落后的消息数据就会丢失。Consumer在消费消息时发生异常导致Consumer端消费失败消费者配置了offset自动提交参数enable.auto.committrue。消费者接受到了消息进行了自动提交。但其实消费者并没有处理完成就宕机了此时kafka认为Consumer已经消费了这条消息了后续便不再分配造成了消息的丢失1.2解决方法——Producer消息发送消息失败关于上面Producer消息发送消息失败的解决方法总结归纳出五种可以结合使用生产者调用异步回调消息。伪代码如下producer.send(msg, callback)生产者增加消息确认机制设置生产者参数acksall。partition的leader接收到消息等待所有的follower副本都同步到了消息之后才认为本次生产者发送消息成功了。生产者设置重试次数。比如retrie3,增加重试次数以保证消息的不丢失定义本地消息日志表定时任务扫描这个表自动补偿做好监控告警。后台提供一个补偿消息的工具可以手工补偿。1.3解决方法——broker写入磁盘失败同步刷盘不太建议。同步刷盘可以提高消息的可靠性防止由于机器没有及时写入磁盘的消息丢失。但是会严重影响性能利用Partition的多副本机制建议。使用下面的这段配置unclean.leader.election.enablefalse表示不允许非ISR中的副本被选举为leader以免数据丢失replication.factor3消息分区的副本个数建议设置大于等于3个min.insync.replicas1这个值大于1要求leader至少能和一个Follower副本保证联系1.4解决方法——Consumer消费异常消费者需要关闭自动提交采用手动提交offsetenable.auto.commitfalse并在代码中写入// 同步提交 consumer.commitSync(); // 异步提交 consumer.commitAsync();代码语言javascriptAI代码解释// 同步提交consumer.commitSync();// 异步提交consumer.commitAsync();2如何避免重复消费消息这实际上是一个消息的幂等性问题幂等性是指一个操作可以被重复执行但结果不会改变的特性。在消息队列中幂等性是指在消息消费过程中保证消息的唯一性不会出现重复消费的情况 。我们有以下几个方案可以解决对于一些业务相关的消息我们通常有需要处理的消息业务主键。比如说发送短信的发送流水号支付业务的订单流水号等。当消费者接受到消息后使用这个消息主键建立获取分布式锁同时将消息业务主键写入库。如果第一步成功消费者进行消费当消费者处理完成后释放分布式锁如果有一条重复的消息进入那么在第一步中就会失败要么是分布式锁要么是数据库主键冲突针对没有业务的消息可以再生产消息的时候给予一个分布式全局ID后面的处理方法与第一条类似在有状态流转的业务当中一个消费者只消费一种业务状态当这个消息的业务状态已经更新、已经处理。那么直接丢弃掉此次消息即可乐观锁消息在生产的时候携带业务上一次查询出的版本号在消费时携带版本号去更新数据库。如果乐观锁原因导致失败那么不需要进行后续处理insert … on duplicate key update消费插入数据时数据已存在则进行更新3kafka的零拷贝是什么原理第一次将磁盘文件读取到操作系统内核缓冲区第二次将内核缓冲区的数据copy 到 application 应用程序的 buffer第三步将 application 应用程序 buffer 中的数据copy 到 socket 网络发送缓冲区(属于操作系统内核的缓冲区)第四次将 socket buffer 的数据copy 到网卡由网卡进行网络传输。如下图取消掉两次CPU的拷贝从而减小CPU的消耗。零拷贝是操作系统提供的如Linux上的sendfile命令是将读到内核空间的数据转到 socket buffer进行网络发送还有Java NIO中的transferTo()方法4kafka如何在分布式的情况下保证顺序消费在kafka的broker中主题下可以设置多个不同的partition而kafka只能保证Partition中的消息时有序的但没法保证不同Partition的消息顺序性。比如说有一个主题Topic A里面有两个Partition但消费端只有一个Consumer。根据上面的架构可以知道这个Consumer会消费两个Partition中的消息这样就肯定会出现消费乱序的情况。那么针对上面这种乱序的情况我们可以这样进行设置一个主题只建立一个Partition这样所有的消息也就只会发送到一个Partition中也就保证了消息的顺序性。Producer也可以指定往一个partition中发送消息。具体可以查看第二章第6节可以保证一个Partition只能被一个Consumer消费也可以保证消息的有序性消费。但也要避免Rebalance原本一对一好好的Consumer宕机或者下线导致Rebalance就会导致消费的乱序。5kafka为什么这么快主要原因有下面几个磁盘写入采用了顺序读写保证了消息的堆积顺序读写磁盘会预读预读即在读取的起始地址连续读取多个页面主要时间花费在了传输时间而这个时间两种读写可以认为是一样的。随机读写因为数据没有在一起将预读浪费掉了。需要多次寻道和旋转延迟。而这个时间可能是传输时间的许多倍。零拷贝第3节提到过避免了两次CPU拷贝减少了CPU的消耗分区、分段、索引再配合二分查找检索提高消息的检索效率分区Partition有效避免了消息的堆积分段segment消息在Partition里面消息是分段来进行存储的每次操作的消息读写都是针对segmengt索引indexkafka分段后的数据建立的索引文件就是第二章第9节的文件存储结构批量压缩读写多条数据一起压缩存储读取kafka是直接操作的page cache而不是堆内对象读写速度更高。且进程重启后缓存也不会丢失6什么是ISR它有什么用在kafka中除了有ISR还有OSRAR功能如下ISRInSyncRepli在kafka中当一个broker宕机挂掉的时候原本在其broker的Leader Partition会重新进行竞选。这个竞选基本从ISR队列中选举。那么现在可以这样说ISR是一个维护了Follower Partition的队列其中的Partition都与Leader Partition消息保持一致。OSROutSyncRepli没在ISR队列中的其他Follower Partition组成的队列ARAllRepli全部分区的Follower Partition也就是ISR和OSR的Partition总和7kafka中的Rebalance是什么什么时候会触发Rebalance是指Partition与Consumer之间的关系需要重新调整分配这个重新调整分配的动作称为Rebalance。那么当出现下面几种情况的时候会触发Rebalance当一个Group中的Consumer新增后当一个Group中的Consumer离开后比如说宕机当Topic下的Partition数量发生变化后总之两边的关系数量发生变化的话都会触发Rebalance8当kafka出现消息积压时该怎么办当出现上面这种情况的时候要么就是Consumer挂掉了或者消费水平太低要么就是Producer消息太多间接导致Consumer消费不及时。针对上面这种情况我们可以有以下的解决方案可以结合使用提高Consumer的数量可以通过增加消费者组中的Consumer数量或者增加Consumer实例来实现。这样每个Consumer可以并行处理消息提高整体消费能力。增加Partition分区数量在kafka中可以设置主题下的Partition将消息分散至更多的Partition中配合第一点方案提高整体的消费能力提高Consumer的消费能力优化消费者的处理能力确保Consumer能够快速处理每条消息。将Consumer处理消息的速度优化至高于Producer生产消息的速度。在不破坏代码业务逻辑的情况下也可以使用异步处理来消费消息。在面试过程中第三点方案是至关重要的很多企业由于硬件资源的原因没有增加Consumer的数量没有增加Partition数量的空间。故此Consumer优秀的消费能力就成了他们考察的目标了。