1. 为什么用问题驱动式来复习Kafka
1.1 学Kafka最容易踩的坑:知识点碎片化
你们有没有这种经历:看了好几篇Kafka教程,知道Topic、Partition、Producer、Consumer这些名词,但真到面试官问“Kafka怎么保证消息不丢失”的时候,脑子里一堆概念呼之欲出,却组织不成一个完整的答案。我这些年带过不少新人,也做过不少技术面试,发现大部分人对Kafka的理解都卡在“名词都会,串联不会”的阶段。
原因很简单。Kafka本身是个分布式系统,知识点天然散布在“生产端、Broker端、消费端”三个不同的领域里,官方文档又写得极简,社区文章则往往只讲某一个点。你单看一篇讲“零拷贝”的文章,觉得懂了;再看一篇讲“幂等性”的,也懂了。但这两者之间有什么关系?它们各自解决的是哪个环节的问题?如果脑子里没有一个完整的地图,知识点就永远是零散的碎片,遇到实际问题根本调不出来。
所以我一直推荐“问题驱动式”的复习方式。不按名词字典的顺序背,而是站在一个使用者的角度,问出一连串真实的、环环相扣的问题——Kafka是什么、怎么保证不丢消息、为什么这么快、消费者是怎么工作的、出了问题怎么排查——每回答一个问题,就把相关的知识点串起来,顺便带出边界情况和常见坑。这种方式学完一遍,知识是形成体系的,不是散落的。
1.2 这篇复习文章适合谁读
这篇是系列的第一篇,聚焦在Kafka的核心架构、消息可靠性、高性能原理和消费端机制上。适合三类人:
- 刚接触Kafka不久,看完基础教程但还没形成知识体系的开发者;
- 准备面试,需要在短时间内把Kafka核心问题梳理清楚的候选人;
- 工作中已在用Kafka,但对底层原理理解不够深、遇到问题时只能靠猜的工程师。
考虑到现在Kafka面试题几乎是Java后端岗位的必考内容,网上相关的信息量又太大太杂,我尽量把这一篇做成“一篇顶十篇”的浓缩笔记:每个问题都给结论,每个结论都给原因,每段原因都尽量配上实际场景。哪怕你之前没怎么用过Kafka,只要按着问题往下读,也能顺畅地跟上节奏。
2. 从“Kafka是什么”出发:核心概念与架构
2.1 Kafka到底解决什么问题
很多人一上来就背“Kafka是一个分布式消息队列”,但这句话其实有点误导。消息队列是它的形态,但它真正解决的是海量数据在异构系统之间的可靠传输与解耦问题。业务系统A产生订单数据,业务系统B需要实时分析,系统C需要归档存储——如果让A直接调B和C的接口,耦合度高不说,B一旦性能抖动,A也跟着遭殃。引入Kafka之后,A只把消息写到Kafka,B和C各自按自己的节奏去消费,互不影响。
这个定位决定了Kafka的设计取向:高吞吐、高可用、数据可回溯。它不像RabbitMQ那样把重点放在灵活的路由规则上,也不像RocketMQ那样在事务消息上花大量功夫,而是把“把大量数据快速、可靠地搬来搬去”这件事做到极致。所以你会看到,Kafka的一个典型特点是没有复杂的路由,Topic就是简单的分类标签。
理解了这个定位,很多后续问题就顺理成章了。比如为什么Kafka的消费模式是拉取而不是推送?因为目标系统消费速度不可控,拉模式才能让消费者自己掌握节奏。再比如为什么Kafka不删除已消费的消息,而是按时间保留?因为数据可回溯本身就是它的设计理念之一。
2.2 Topic、Partition和Offset是绕不开的三兄弟
这三个概念是Kafka一切机制的地基,必须彻底吃透。
**Topic(主题)**是逻辑上的消息分类。比如一个电商系统,可以有order_topic、payment_topic、user_topic,每个Topic里的消息类型不同。Topic本身不存储数据,它只是一个逻辑名字。
**Partition(分区)**是物理上的存储单元。每个Topic可以被划分为多个Partition,每个Partition是一个有序的、不可变的日志文件(log)。消息写入Partition时是追加写,读取时按offset顺序读。为什么要分区?一句话:为了并行。单台机器无论怎么优化,磁盘写入总有上限,把Topic拆成多个Partition分布到多台Broker上,就能同时写多个磁盘,吞吐量成倍提升。
但是分区也带来了副作用——跨分区的消息顺序无法保证。假设用户下单的消息被路由到Partition 0,支付消息被路由到Partition 1,消费端看到的顺序就可能错乱:先看到支付成功,再看到下单。所以Kafka的顺序保证是“分区内有序”,不是“全局有序”。如果你严格要求全局顺序,就只能把Topic设置成1个分区,但这就牺牲了吞吐量。这是一对需要业务层面做权衡的矛盾。
**Offset(偏移量)**是消息在Partition内的位置编号,从0开始,单调递增。每一条消息在所属Partition内都有一个唯一的offset,相当于它在日志文件里的游标。消费者消费的时候,本质上就是在每个Partition上维护一个“当前读到哪了”的指针,这个指针就是offset。
理解offset还有一个关键点:Kafka不会因为消息被消费就删除它。消息会按照broker的保留策略(默认7天或根据大小)保留一段时间,消费者可以按需重置offset,从头再读一遍,或者跳到某个时间点。这个设计给“数据回溯”和“重建消费逻辑”提供了极大便利,是我个人认为Kafka最被低估的优势之一。
2.3 Broker、Producer、Consumer怎么配合
Kafka集群里跑着若干个Broker节点,每个Broker负责一部分Partition的读写。Producer发送消息到指定Topic,由分区器(Partitioner)决定消息进哪个Partition;Consumer订阅Topic,从各个Partition拉取消息进行处理。
这里有个经常被忽略的角色——Controller Broker。在Kafka集群中,会有一个Broker被选举为Controller,负责管理整个集群的元数据变化,比如Partition的Leader选举、Broker的上线与下线等。很早期的Kafka版本依赖ZooKeeper来做这件事,现在Kafka 2.8之后引入了KRaft模式,逐步去掉ZooKeeper依赖。从3.x版本开始,KRaft已经进入生产可用状态。如果你在2024年之后新装集群,可以直接考虑KRaft模式,少维护一套ZooKeeper系统,省心不少。
Producer和Consumer之间不直接通信,所有消息都经过Broker中转。这种架构天然做到了生产和消费的解耦:Producer不需要知道Consumer在哪,Consumer也不需要知道Producer是谁。
还有一个概念容易被漏掉:Consumer Group(消费组)。同一个Group里的多个Consumer共同消费一个Topic时,每个Partition只会被组内一个Consumer消费,目的是水平扩展消费能力。关于消费组,后面第5节会重点展开,这里先留个钩子。
3. 消息可靠性:Kafka怎么做到不丢消息
3.1 Producer端的确认机制怎么选
“Kafka到底会不会丢消息”这个问题,面试被问到的频率极高,回答时一定要分三段来谈:生产端、Broker端、消费端。任何一段没做好,都可能丢消息。
生产端最关键的是acks参数。Producer发消息给Broker时,需要Broker确认“我收到了”,acks控制的就是这个确认的严格程度:
| acks配置 | 行为 | 可靠性 | 适用场景 |
|---|---|---|---|
| 0 | 发出去就不管了,不等待任何确认 | 最差,可能丢消息 | 日志、监控等可容忍丢失的场景 |
| 1 | Leader Partition写入成功即确认 | 中等,Leader宕机时可能丢 | 大多数业务场景 |
| -1 或 all | 所有ISR副本都写入成功才确认 | 最高,最安全 | 金融、订单等关键数据 |
我用“领导签字”来类比说明。acks=0相当于把文件往领导桌上一扔就走,领导有没有看你不知道;acks=1相当于等领导签了个字,但领导告诉你“没问题”的时候,他可能还没让其他同事复核过;acks=all相当于领导不仅自己签了字,还等所有参与复审的同事都签完字才告诉你“办妥了”。自然是越往后越稳,但耗时也越长。
实际业务中我建议:只要能接受一点延迟,就用acks=all + 开启幂等。很多人觉得acks=all会拖垮吞吐,实测下来并没有那么夸张,配合批量发送和压缩,吞吐损失通常在可接受范围内。而一旦数据出问题,丢消息的代价远大于这点延迟。
3.2 副本同步与ISR机制
Broker端保证不丢消息靠的是副本(Replica)机制。每个Partition可以配置多个副本,其中一个是Leader,负责读写;其余是Follower,只负责同步Leader的数据。为什么要Follower?如果Leader所在的Broker宕机了,得有一个Follower顶上,这就是高可用的基础。
Follower是主动从Leader拉数据,拉取快慢取决于网络和自身负载。为了维护“哪些副本是健康的”这个状态,Kafka定义了ISR(In-Sync Replicas,在同步副本集合)。ISR里的副本必须满足两个条件:与Leader之间副本数差距在可接受范围内,且在指定时间间隔内与Leader保持心跳。
这里有一个关键点:acks=all等待的是ISR里的副本全部确认,不是所有副本。如果一个Follower落后太多,被踢出ISR,那acks=all就只等当前ISR成员确认。这保证了“消息不丢”和“可用性”之间的平衡——如果要求所有副本都确认,那只要一个副本挂了,整个分区就写不进去了。
这里我建议记住一个黄金组合,面试或实际配置都可以直接说:
acks=allmin.insync.replicas=2unclean.leader.election.enable=false
min.insync.replicas=2的意思是:如果ISR里少于2个副本,Broker直接拒绝写入,宁可不可用,也绝不丢数据。这个配置在金融级场景几乎是标配。unclean.leader.election.enable=false则禁止那些落后于Leader的副本参与Leader选举,防止“选出来一个新的Leader,但它丢了很多数据”。
3.3 消费端到底怎么才算“消费成功”
很多人在生产端、Broker端做了大量配置,却忽略了消费端同样能丢消息。简单说:如果你用的是自动提交offset,那在消息处理成功之前,offset就可能已经提交了。
Kafka消费者有一个enable.auto.commit配置,默认是true,每auto.commit.interval.ms(默认5秒)自动提交一次当前消费到的offset。问题在于:假设你poll了一批消息,正在处理第3条时突然宕机重启,offset在上一次自动提交时已经往前走了,重启后Consumer会从上次提交的位置继续读,那3条消息就永远不处理了。
解决办法:把enable.auto.commit设为false,改为手动提交offset。具体时机很重要——一定要等所有消息都处理成功后再提交,不能一边处理一边提交。使用Kafka Consumer的Java API时,代码大致是这样:
while (running) { ConsumerRecords<String, String> records = consumer.poll(Duration.ofMillis(1000)); for (ConsumerRecord<String, String> record : records) { // 处理业务逻辑 process(record); } // 全部处理成功后才提交 consumer.commitSync(); }如果你用的是commitSync,它是同步阻塞的,提交失败会抛出异常并重试,保证“要么提交成功,要么抛异常”;如果你用commitAsync,它是异步的,不阻塞,吞吐更好,但失败时没有重试机制,极端情况下可能提交失败。我的建议是:关键场景用commitSync,或在commitAsync的回调里做好失败补偿。
这类问题是“消息重复消费”和“消息丢失”之间反复横跳的根源。真正做到不丢消息后,你会遇到另一个问题——消息可能重复。比如消息处理成功但提交offset前宕机了,重启后就会再次消费到这条消息。所以对于需要精确处理的消息,业务侧必须做幂等性设计,例如在数据库里用唯一索引约束,或用Redis Set记录已处理的key。这是分布式系统里没办法绕开的现实。
4. Kafka高性能的底层逻辑
4.1 为什么Kafka能把吞吐量做到百万级
单机Kafka就能做到每秒百万条消息的写入量,这个数字在消息队列里属于天花板级别。很多人觉得不可思议,但拆开看其实每个环节都是“常识的组合”。
第一层是利用磁盘顺序写的特性。磁盘的随机写确实慢,但顺序写完全不一样——机械硬盘的顺序写可以跑到每秒一两百MB,SSD更是能轻松上GB。Kafka的每个Partition就是一个追加写的日志文件,消息永远只追加到文件末尾,不存在随机寻道的问题。这就相当于一个仓库管理员,每次进货都往货架最后面摆,速度当然快。如果你让他在仓库里随机找空位摆,就得反复走来走去,效率就下来了。
第二层是Partition并行。单条日志流再快也有瓶颈,Kafka把Topic拆成多个Partition分散到多台Broker上,相当于多个仓库管理员同时干活。集群整体吞吐量约等于“单分区吞吐量 × 分区数”,这也是为什么当你发现Kafka集群吞吐上不去时,第一反应往往是“分区数够不够”。
第三层是网络和磁盘的异步化。Kafka的Producer端可以把多条消息攒一批再发,Broker端攒一批再落盘,Consumer端攒一批再回调处理函数。整个链路上“攒批”的动作无处不在,极大摊薄了每一条消息的系统调用开销。
4.2 顺序写和页缓存是如何配合的
顺序写解决了磁盘写入速度问题,但Kafka的读取也很快,尤其在“刚写入的消息马上被消费”这种典型场景下。这就轮到**页缓存(Page Cache)**登场了。
Kafka并不自己管理消息缓存,而是直接依赖操作系统的Page Cache。写入消息时,数据先拷贝到内核Page Cache里,由操作系统异步刷到磁盘;读取消息时,优先从Page Cache里读,只有Page Cache没命中才走磁盘。如果你的消费速度跟生产速度差不多,那大部分读取都会直接命中Page Cache,根本碰不到磁盘。
这种设计的高明之处在于两点。第一,省掉了JVM堆内缓存带来的GC压力——如果Kafka自己用堆内存缓存大量消息,JVM的垃圾回收会频繁停顿,吞吐必然受拖累。第二,操作系统对Page Cache的管理已经足够优秀——它知道哪些页是热的,哪些可以淘汰,比应用自己管理缓存更高效。
实际调优时,你可以观察到Kafka进程的物理内存占用往往看起来很低,因为它把大量空闲内存交给了Page Cache。如果你想验证效果,可以做一个简单测试:同步生产一条消息后立刻消费,几乎感觉不到延迟,因为数据还在Page Cache里。但如果你把消费暂停很久再进行消费,就会发现读取耗时明显上升——因为数据早已被刷到磁盘,Page Cache也早就淘汰了这些页。
4.3 零拷贝到底省了什么
“零拷贝”是Kafka面试中的明星知识点,但对它的理解经常停留在“不用拷贝所以快”的层面,没有真正说清楚省的是什么。我来拆解一下。
传统情况下,从磁盘读取文件并通过网络发送给客户端,数据要经历四次拷贝:磁盘→内核缓冲区→应用缓冲区→Socket缓冲区→网卡。每次拷贝都有CPU开销和上下文切换开销。
Kafka在消费者读取消息时用了sendfile系统调用,数据直接从内核的文件缓冲区发送到网卡,跳过了应用层缓冲区。也就是从四次拷贝减少到两次,内核态和用户态的切换次数也大幅减少。对于日志类文件这种“整块数据直出”的场景,效果极其显著。
为什么Kafka能做零拷贝而很多其它系统做不了?因为Kafka的消息是不可变的、按顺序存储的日志文件,消费时只需要按offset定位文件区域,然后把这一整块数据原样发出去,不需要任何应用层的加工——天然适配零拷贝。而像数据库那种读取后要做权限校验、行过滤、格式转换的场景,想做也做不了。
我还见过一些新手把Java的MappedByteBuffer(内存映射)和零拷贝搞混。两者都能减少拷贝,但MappedByteBuffer是把文件映射到进程地址空间,依然会发生用户态和内核态之间的切换;而sendfile则完全在内核态完成数据发送。理解这点区别,面试时能加不少分。
5. 消费端核心机制与常见坑
5.1 推还是拉?Kafka为什么选拉模式
几乎所有消息队列的消费模型都是“Broker推给Consumer”或“Consumer从Broker拉”二选一。Kafka选择了拉模式(pull),原因很直接:推模式没法处理消费速度不均的问题。
如果Broker主动推,它不知道消费者当前能不能处理得过来。推快了,消费者积压;推慢了,消费者闲置。而且推模式下,消费者需要维护一个复杂的窗口状态来反馈“我还能收多少”,Broker也需要为每个消费者维护发送状态,复杂度直线上升。
拉模式的逻辑就简单多了:消费者按照自己的节奏,主动去Broker上取数据,一次取多少、多久取一次,完全自己决定。代码里常见的就是poll(Duration.ofMillis(1000)),Consumer在循环里不断调用poll,Broker每次返回一批可消费的消息。
拉模式有一个副作用:如果Topic里暂时没有新消息,poll会一直空转,白白消耗CPU。Kafka的处理方式是poll(timeout)让线程阻塞指定时间后才返回空结果。这个超时时间值得留意——设置过短会导致CPU空转严重,设置过长会导致消息处理不实时。我个人常用的起始值是500到1000毫秒,再根据业务实时性要求调整。
5.2 Consumer Group与Rebalance机制
前面提过,同一个Consumer Group里的消费者共同消费一个Topic时,每个Partition只会分配给组内的一个Consumer。这样设计的目的是水平扩展消费能力:如果你只有一个消费者,那它要处理所有分区的数据;如果你把消费者数量增加到3个,它们就能并行消费。
这里有一个必须记住的规则:一个Group内的消费者数量超过分区总数时,多出来的消费者会空闲。很多新手以为消费者越多消费越快,结果把消费实例开到比分区数还多,性能完全没有提升,反而多带来了一堆空转线程。比如某个Topic有10个分区,你最多只能开10个消费者并行消费,第11个开了也拿不到任何分区分配。
当Group里的消费者数量发生变化,或订阅的Topic分区数发生变化时,Kafka会触发Rebalance(再平衡):把所有分区收回,重新分配给组内消费者。这个概念本身不难,但坑在于Rebalance期间,整个Group是暂停消费的。如果你频繁重启消费者,或者某个消费者处理消息太慢导致超出了max.poll.interval.ms(默认5分钟),就会反复触发Rebalance,造成“消费卡死”的恶性循环。
我踩过的一个坑是:消费者里有段代码偶尔会长时间阻塞(比如调了一个慢SQL),超过5分钟没调用poll,Consumer就被判定为“死了”,触发Rebalance把它的分区分给别人。等它缓过来,又要再触发一次Rebalance把分区要回来。整个过程就像一个办公室里同事频繁请假,项目不停换人,进度反而更慢。解决办法是:把消费和处理分离,poll之后把消息放入线程池异步处理,保证poll循环不被阻塞。
5.3 Offset提交:手动提交还是自动提交
面试时这个问题也高频出现:“你平时用的是自动提交还是手动提交offset?”如果回答“用的默认自动提交”,面试官大概率会继续追问“那会不会丢消息”。正如前面3.3节提到的,自动提交存在消息未处理完成就提交offset的风险。
这里想把三种提交方式放在一起对比:
| 提交方式 | 配置/API | 优点 | 缺点 |
|---|---|---|---|
| 自动提交 | enable.auto.commit=true | 省心,代码简单 | 可能丢消息 |
| 手动同步提交 | commitSync | 可靠,失败会重试 | 可能拖慢消费速度 |
| 手动异步提交 | commitAsync | 不阻塞,吞吐高 | 失败无重试,可能丢offset |
如果业务对吞吐要求高,我的习惯做法是commitAsync+回调里记录失败日志并做补偿。如果数据绝对不能丢,就用commitSync。还有一种场景很特殊:你希望在“数据刚好处理完”时提交,可以用commitSync(Map<TopicPartition, OffsetAndMetadata>)精确指定要提交的offset值,这在实现“至少一次”语义时很常用。
再补充一个关于offset的冷门操作:设置auto.offset.reset为earliest或latest。当消费者没有已提交的offset时(比如新消费者首次启动,或offset因过期而被删除),这个配置决定它从哪开始读——earliest是从最早可用的消息开始读,latest是只读新消息。这个参数看似不起眼,实际踩坑率很高,我就见过有同学在生产环境新建一个消费者,期望它“从当前时间开始消费后续消息”,结果忘了设置latest,默认的latest其实没问题,但如果你期望从最早开始、却发现消费者跳过了历史数据,那就要检查这个参数是否设错了。
6. 高频异常排查实践
6.1 消息延迟高怎么定位
很多同学在生产环境看到“Kafka消息延迟高”就慌了,上来就重启消费者,但这样做不一定解决问题,甚至可能让情况更糟。我一般按下面的顺序排查:
第一步,先确认延迟到底发生在哪段链路。是Producer发消息慢?还是消息到了Broker但消费者不消费?还是消费者在处理消息时耗时高?最简单的方法是看消费组的Lag(堆积量)——如果Lag持续增长,说明消费速度跟不上生产速度;如果Lag不高但业务感觉消息“来得慢”,那问题可能出在Producer端或网络链路。
第二步,如果确认是消费端跟进不上,看消费者实例的数量和分区的数量关系。消费者数是否等于分区数?如果小于,恭喜你找到了瓶颈——增加消费者或线程就能提升吞吐。如果已经等于分区数,就得往下看单条消息的处理耗时,比如Consumer的日志里是否频繁出现单条消息处理耗时几百毫秒以上,那就是业务代码的锅了。
第三步,检查fetch.min.bytes和fetch.max.wait.ms参数。这两个参数控制着Consumer拉取消息的批量大小和等待时间。默认情况下,Consumer只要拉到少量数据就会返回,但如果你把fetch.min.bytes设得太高,Kafka就会一直等数据攒够了才返回,即使新消息早就到了,Consumer端也会感到“延迟”很大。这类延迟是人为配置出来的,改小即可。
6.2 消费者堆积是常态,怎么处理
消费者堆积(Lag持续增长)在生产环境几乎是无法避免的,大促期间尤其常见。处理思路要分短期和长期来看。
短期的应急手段是快速扩容消费者数量。但这里有个前提:如果Topic分区数没有变多,多开的消费者实例只能干瞪眼。所以为了留出扩容空间,我建议在创建Topic时就把分区数设置得比预估的峰值消费并发大一些,比如预估需要5个消费者并行,那分区数至少设10个以上。记住一个公式:最大消费并行度 = 分区数,分区数从一开始就要留好余量。
长期来看,堆积的根因往往是消费速度跟不上生产速度。常见解法包括:对消费逻辑做批处理优化,比如一条一条更新数据库改成批量更新;调整max.poll.records来控制每次poll的消息数量,避免单次拉取太多导致处理时间过长而触发Rebalance。
特别提醒一个容易忽略的点:Kafka的Lag是从Consumer Group的offset提交情况算出来的,如果某个消费者一直在处理但从不提交offset,Lag会越来越大。这时你看到的“堆积”其实是假象,实际是消费端在重复处理或处理逻辑卡住了。先看一眼Consumer日志有没有在正常输出,再决定要不要扩容。
6.3 集群宕机后的恢复顺序
Kafka集群宕机分很多种:单台Broker宕机、整个机房网络故障、元数据故障等,恢复策略各不相同。我经历过最典型的是单台Broker宕机后Kafka集群不可写,排查下来发现是配置问题。
如果Topic的副本因子是3,某台Broker挂了,只要还有一台ISR副本活着,理论上集群应该自动完成Leader切换,继续对外服务。但如果你的min.insync.replicas=3且ISR只剩2个副本,那么所有写入请求都会被拒绝,表现为“Kafka突然写不进去了”。这时候要判断:是保护数据宁可不可用,还是牺牲数据保可用性?我的建议是,业务关键数据必须保数据,但要做好告警和预案,而不是生产环境遇到问题了才临时决策。
还有一类宕机是Controller挂了或Controller与其它Broker失去联系。Controller负责选举Leader和生成元数据变更,它挂了会触发一次新的Controller选举,期间集群的元数据操作(比如创建Topic、分区Leader变更)会短暂不可用,但已有的数据读写不受影响。所以如果你发现Kafka集群“部分功能不可用”,先确认Controller的状态。
恢复顺序上,我个人的经验是:先看Broker进程是否还活着,再看集群中是否有足够的ISR副本在线,最后检查所有Topic的Leader是否都已恢复。Kafka自带的kafka-topics.sh --describe命令可以一次列出所有Topic的Leader、ISR和副本分布,排障时几乎是必查项。
7. 面试与实操的高频考点整理
7.1 面试题里的Kafka高频问题
这一节帮大家把面试中围绕Kafka最常被问的问题整理成一个速查表。每个问题都尽量给到“一句话结论+扩展要点”,方便背诵和理解:
| 问题 | 一句话结论 | 扩展要点 |
|---|---|---|
| Kafka是什么? | 分布式消息队列,解决数据缓冲、解耦和异步问题 | 高吞吐、高可用、可回溯 |
| 消息顺序怎么保证? | 分区内有序,跨分区不保证 | 全局有序需单分区,牺牲吞吐 |
| 怎么保证消息不丢? | 生产端acks=all,Broker端副本同步,消费端手动提交offset | 三端配合,缺一不可 |
| 怎么避免消息重复? | 做不到完全避免,需要业务幂等 | 唯一索引、Redis去重 |
| Kafka为什么快? | 顺序写、Page Cache、零拷贝、批量传输 | 四大原因要能展开讲细 |
| 消费者和分区是什么关系? | 一个分区在同一时刻只能被group内一个消费者消费 | 消费者数不能超过分区数 |
| Rebalance是什么? | 消费者增减或分区变化时,分区重新分配的过程 | Rebalance期间暂停消费,要避免频繁触发 |
| 推还是拉? | 拉模式 | 消费者自己控制节奏 |
| 如何提高消费吞吐? | 增加分区数、增加消费者、批量处理 | 消费者数要匹配分区数 |
| Leader选举逻辑? | 从ISR中选,尽量不选非同步副本 | unclean.leader.election.enable踩坑 |
这类问题其实是上文所有内容的一个浓缩。如果你能把每一行的“扩展要点”讲给面试官听,基本能覆盖Kafka大部分考点。
7.2 实操中的5个避坑技巧
最后分享几个我这些年实际踩过、看过别人踩过的坑,每一个都对应着具体的规则:
坑一:Topic分区数拍脑袋决定,后期无法缩小去改。
分区数只能增加,不能减少(除非删掉重建Topic)。创建Topic时如果暂时无法准确预估流量,宁可多分一些。分区太多也会带来副作用,比如文件句柄占用多、Rebalance耗时长,所以建议按“未来峰值消费并发数×1.5”左右的粗量来创建。
坑二:消费者max.poll.interval.ms设置不当引发连环Rebalance。
如果你的单条消息处理时间有可能超过默认5分钟,一定要显式调大这个参数。同时把处理逻辑独立成异步线程,保持poll循环活跃,是最稳妥的做法。
坑三:忽略client.id的作用。分布式环境下,每个Consumer最好设置一个有辨识度的client.id,比如${appName}-${hostname}-${random}。这不仅让监控面板清晰可见,排障时也能快速锁定这台实例是哪个环境的哪个机器。
坑四:对Kafka的监控只关注机器CPU、内存,不关注Lag。
Kafka集群健康的核心指标其实在消费端的Lag。建议每台消费机器上都配上Lag监控,Lag持续上涨时要能触发告警。用Kafka官方自带的kafka-consumer-groups.sh --describe也能手动查看,但生产环境最好交给监控系统。
坑五:盲目使用acks=0追求高吞吐。
很多教程会说“日志类业务可以用acks=0”,但这句话的前提是你真的能接受丢失。我见过有团队把订单入库的消息也配了acks=0,结果某次Broker抖动直接丢了几万条消息,事后排查悔之晚矣。我的底线是:凡是能对账、对用户有直接影响的数据,一律不用acks=0。
结尾
这套问题驱动式的回顾,第一遍读下来大概需要半小时到一小时。花这段时间是值得的,因为你收获的是一张完整的Kafka知识地图。我自己复习Kafka时最喜欢用的方式,就是假装自己是面试官,把每个知识点周围的“为什么”全部问一遍,再假装自己是排查故障的运维,把每个故障场景从头走一遍。两轮走完,大部分问题都能在脑子里直接形成清晰的答案。
下一篇我会继续沿着问题驱动的方式,聊Kafka的Producer和Consumer API细节、消息压缩与序列化、Kafka Streams、MirrorMaker跨集群同步等更进阶的内容。如果你读完这篇有觉得不清楚或者和自己理解不一致的地方,欢迎在评论区留言讨论,我看到了会尽量回复。