
我第一次认认真真研究 Kafka 消息队列不是赶时髦而是因为凌晨两点的线上故障订单接口超时被打爆库存扣减、积分发放、物流通知全挤在一次请求里同步调用任何一个下游服务抖动用户就要陪着一起等。折腾到天亮我意识到问题不是出在代码写得不够快而是架构里缺了一个组件——它能让系统在高峰时先接单、后处理把突发的流量削平把服务之间的强耦合拆开。这个组件就是 Kafka。这篇文章不是把官方文档翻译一遍而是站在我实际用 Kafka 做生产项目的角度从零开始讲清楚消息队列是什么、Kafka 怎么装、生产者和消费者怎么写、线上重复消费和延迟高怎么排查最后再把面试里最容易翻车的几个原理题串一遍。适合完全没接触过 Kafka 的后端开发也适合已经会发消息但被各种坑折磨过的同学。1. 消息队列到底在解决什么问题1.1 同步调用的痛我从一个订单接口说起先想象一个非常常见的下单流程用户点支付订单系统需要同步调用库存系统扣库存、调用优惠券系统核销、调用积分系统加积分可能还要推一条短信。这个接口的总耗时等于所有下游耗时之和。有一天库存系统发版慢了两百毫秒整个下单接口集体变慢再有一天积分系统数据库连接池被打满下单接口直接超时。用户不会管是不是积分系统的锅他只看到下单失败。这就是同步架构的问题——任何一环不稳定整个链路都会被拖下水而且系统越滚越大谁也不敢轻易加新的下游服务因为每加一个同步调用请求就要多等一次网络往返。1.2 解耦、异步、削峰消息队列被发明出来的理由消息队列的解法很直接把同步调用变成异步消息。订单服务把订单已创建这个事件写进 Kafka库存、积分、短信各自从 Kafka 里订阅这个事件拿到之后自己处理。这样一来有三大好处异步化下单接口不用等库存和积分处理完只要消息写进 Kafka 就算成功响应时间从几百毫秒降到几十毫秒。解耦订单服务完全不需要知道下游有哪些系统。以后新增一个数据分析服务直接订阅同一条消息就行订单服务一行代码都不用改。削峰填谷秒杀场景下流量瞬间暴涨数据库一瞬间接不住。Kafka 先把海量请求按自己的节奏接进来下游消费服务再按自己的处理能力批量跑高峰期不崩高峰期过后慢慢追平。我经常拿餐厅打比方没有传菜窗口时厨师炒好一盘菜得满大厅找对应的服务员炒一个催一个最后锅铲都要抡出火星。有了传菜窗口厨师只管把菜放进窗口、按铃服务员忙完手头的活再过来取。窗口能积压多少菜决定系统能扛多大流量——这就是消息队列的缓冲价值。1.3 Kafka 和其他消息队列怎么选消息队列不是一个新概念RabbitMQ、RocketMQ、ActiveMQ、甚至 Windows 上的 MSMQ 都在各种遗留系统里服役。但 Kafka 的定位和它们不太一样对比维度KafkaRabbitMQRocketMQ核心设计分布式日志追加写AMQP 协议多交换机路由金融级消息事务丰富吞吐量极高百万级每秒中等几万级每秒高几十万级每秒消息顺序分区内有序需绑定队列与消费者分区内有序消息堆积能力强可长期堆积堆积能力一般较强最典型场景日志、大数据、事件驱动中小系统业务解耦电商交易、金融对账Kafka 最初是 LinkedIn 为了处理海量日志和用户行为数据开发的所以底层设计天生向吞吐量和堆积能力倾斜。它把每条消息当作日志里的一个记录不断追加写入磁盘所以非常抗压。如果在选型时遇到哪种场景适合 Kafka我的判断标准就一句话当你需要高吞吐、可回放、能堆大量消息的时候Kafka 合适如果你只是业务量不大、希望路由规则灵活RabbitMQ 学起来更快运维成本也更低。1.4 什么场景不该上 Kafka我见过不少团队把 Kafka 当万能药什么业务都往里塞最后反而被复杂化。如果你的调用链路本来只有两三个服务同步调用的响应时间也能接受那就别引入 Kafka——一个消息队列会带来消息重复、乱序、积压、消费失败重试等一堆新问题。另外强事务要求的场景也要谨慎Kafka 虽然有事务 API但它的定位是最终一致不是像数据库那样提供强一致的事务保证。比如扣款和加积分绝对不能发了个消息就算成功必须设计对账和补偿机制。2. 从零搭一个能跑的 KafkaWindows、Linux 都别慌2.1 运行前准备JDK、版本与 KRaft 模式Kafka 是 JVM 系的消息队列所以装之前先确认机器上有 JDK。Kafka 3.x 之后对 Java 版本要求是 11 及以上部分新版本已经需要 Java 17。如果你手头正好有多个 JDK 版本建议给 Kafka 单独指定JAVA_HOME避免项目里的 Java 8 环境把它带崩。还有一个让新手绕圈子的坑是 ZooKeeper。旧版 Kafka 必须依赖 ZooKeeper 存元数据、做选举搭个集群等于同时部署两套系统非常费劲。Kafka 3.3 开始引入 KRaft 模式把元数据管理收回到 Kafka 自身单节点或三节点集群都可以不用 ZooKeeper。我建议 3.x 的新项目直接用 KRaft 模式少一个组件就少一半故障点。如果你看老文章还在教 ZK 那一套可以先确认一下版本不是文章错了是时代变了。2.2 单机版 20 分钟跑通先去 Kafka 官网下载二进制包文件名类似kafka_2.13-3.7.0.tgz其中2.13是编译用的 Scala 版本3.7.0是 Kafka 版本初学者不用太纠结。下载后解压tar -xzf kafka_2.13-3.7.0.tgz cd kafka_2.13-3.7.0KRaft 模式第一次启动需要先格式化存储目录。这步很多新手会漏漏了就报Storage directory ... not formatted。正确姿势是这样# 生成一个集群 ID KAFKA_CLUSTER_ID$(bin/kafka-storage.sh random-uuid) # 格式化日志目录 bin/kafka-storage.sh format -t $KAFKA_CLUSTER_ID -c config/kraft/server.properties # 启动服务 bin/kafka-server-start.sh config/kraft/server.properties看到started (kafka.server.KafkaRaftServer)之类的日志说明服务起来了。然后开另一个终端验证bin/kafka-topics.sh --bootstrap-server localhost:9092 --create --topic my-topic --partitions 3 --replication-factor 1 bin/kafka-console-producer.sh --bootstrap-server localhost:9092 --topic my-topic bin/kafka-console-consumer.sh --bootstrap-server localhost:9092 --topic my-topic --from-beginning在 producer 窗口敲几行字consumer 窗口能看到说明通路已经打通。2.3 Windows 上的安装要点Windows 上跑 Kafka 也是可行的我甚至见过不少同学的项目是 Windows 开发机连远程测试集群本地只装一个单机版做实验。流程基本一样只是把命令行换成.bat:: 设置 JAVA_HOME 后执行 bin\windows\kafka-server-start.bat config\kraft\server.properties bin\windows\kafka-topics.bat --bootstrap-server localhost:9092 --create --topic my-topic --partitions 3 bin\windows\kafka-console-producer.bat --bootstrap-server localhost:9092 --topic my-topic这里要特别提醒三个 Windows 环境下的坑解压路径不要带中文、空格和特殊符号Kafka 对路径的处理在 Windows 上很娇气放在C:\Program Files下可能直接起不来。如果日志抛出OutOfMemoryError检查KAFKA_HEAP_OPTS或KAFKA_JVM_PERFORMANCE_OPTS默认堆内存设置可能不适用于小机器可以显式设成-Xmx512m -Xms256m起步。Windows 自带杀毒软件有时会扫描日志目录导致 Kafka 写入卡顿可以在测试环境把 Kafka 目录加入白名单。2.4 server.properties 里最影响实战的几个配置不管用 KRaft 还是 ZooKeeper 模式最终都要面对config/kraft/server.properties旧版是config/server.properties。这几个配置我建集群时几乎必改配置项作用我的建议listeners服务对外监听的地址单机学习用PLAINTEXT://localhost:9092跨机器访问要改成内网 IPlog.dirs消息数据落地目录不要放在系统盘数据量一大就很被动num.partitions自动创建主题时的默认分区数默认 1生产按流量设成 3~12log.retention.hours消息保留时长默认 168 小时7天日志型主题可以更长message.max.bytes单条消息最大字节数默认约 1MB大消息场景需要调大这里有个很容易踩的误区很多人以为分区数越多吞吐越高于是建主题直接设 64 个分区。分区数确实能提升并发度但每个分区在消费者、副本、索引层面都有额外开销分区太多而机器太弱反而把性能拉垮。分区数应该围绕目标吞吐估算单个分区每秒处理量 × 分区数 ≥ 峰值消息量再留 30% 冗余。比如单分区实测能扛 500 条/秒业务峰值每秒 1000 条分区数给 4 个就够。我的建议是学习阶段老老实实用默认 3 分区先把整条链路跑熟再去纠结分区数的微调。3. 生产者把消息发得又快又稳3.1 主题、分区、副本三个名词一次搞清楚Kafka 里最基础的概念是主题、分区和副本。主题是消息的逻辑分类好比一本《订单流水账》。分区是主题在物理上的拆分好比这本账被拆成好几册每册存一部分订单。消息并不是一股脑写进主题而是先被路由到某个分区再往这个分区的文件末尾追加。副本则是分区的备份。每个分区默认可以配置多个副本其中一个叫 leader其他人叫 follower。生产者和消费者只跟 leader 交互follower 默默同步 leader 的数据一旦 leader 挂了再从中选一个接管。这三个概念理解了后面所有操作都有抓手。3.2 路由规则为什么 key 如此重要生产者把一个消息发出去核心代码只有一行props.put(bootstrap.servers, localhost:9092); props.put(key.serializer, org.apache.kafka.common.serialization.StringSerializer); props.put(value.serializer, org.apache.kafka.common.serialization.StringSerializer); KafkaProducerString, String producer new KafkaProducer(props); // 带 key 的消息 producer.send(new ProducerRecord(order-events, order-123456, ORDER_CREATED));第二个参数key非常关键。Kafka 默认分区器处理 key 的逻辑是对 key 做哈希再用分区数取模决定消息去哪个分区。于是同一个 key比如同一个订单号的所有消息都会被送到同一个分区也就天然保证了分区内的顺序性。如果 key 传了 nullKafka 会用粘性分区策略在多个分区之间轮询分发这样并发写性能更好但没有顺序保证。所以业务上如果需要同一用户的操作必须按顺序处理就必须带上用户 ID 作为 key而不是图省事传 null。3.3 生产者参数调优克制比堆参数重要生产者的性能瓶颈通常不在 Kafka 端而在生产者自己怎么攒批、怎么确认、怎么压缩。几个核心参数表格列出来参数默认值作用调优方向acksallleader 等多少个副本确认后才算发送成功0 最快但可能丢all 最稳linger.ms0每条消息等多久再批量发送稍微调到 5~20 ms吞吐会明显提升batch.size16 KB每个批次最多攒多少字节大 batch 能提高吞吐但延迟上升buffer.memory32 MB生产者缓存未发送消息的内存大量积压时调大compression.typenone消息压缩方式生产建议 lz4 或 zstd压缩率越高网络压力越小初学者最容易犯的错是acks0图快。在本地 demo 里感觉不到差别生产环境一旦 broker 抖动消息就悄悄丢了而且无迹可查。我一般建议默认值 all 不要动除非你明确知道自己要牺牲多少可靠性换取多少性能。还有一个容易被忽略的是单条消息大小。Kafka 默认单条消息上限约 1MB如果业务要传大 JSON 或文件流需要联动修改三个地方broker 端message.max.bytes、生产者max.request.size、消费者fetch.max.bytes。只改一处是没用的这是接收 1M 消息失败类问题最常见的病因。3.4 生产端写消息的习惯比 API 重要调参只是锦上添花真正影响稳定性的写消息习惯往往被忽略。我总结了三条不要在 for 循环里同步等发送结果。一次send()是异步的如果你紧接着用producer.flush()或直接拿get()等结果等于把异步又改回同步高吞吐场景立刻打回原形。正确做法是发送回调里记录失败日志定时flush()。失败重试要有界限。retries默认已经不小但网络长时间抖动时无限重试会拖死生产者。配合delivery.timeout.ms设定整体上限超过上限就进死信流程别跟一条消息死磕。幂等生产者默认就开。Kafka 3.0 之后enable.idempotence默认是 true配合acksall可以避免重试带来的重复消息。如果在老版本环境记得显式开启。4. 消费者重复消费、消费组和顺序性问题一次说清4.1 消费组和分区分配谁消费哪个分区谁说了算消费者不是孤零零自己跑的它一定要归属于某个消费组也就是配置里的group.id。同一个消费组内的多个消费者会平分主题里的所有分区。比如一个主题有 6 个分区消费组里有 3 个消费者那么每个人消费 2 个分区如果组里只有 1 个消费者它一个人消费全部 6 个分区。这里经常有人误解以为启动两个消费者同一个消息就能被两个服务都收到。不对只有不同group.id的两个消费者组才会各自收到一份完整的消息。这就像订单事件既发给了库存组又发给了风控组两个组各收一份而同一个组内的多个消费者只是同一份消息在不同分区上的分工。分区分配不是固定的。当消费者加入、离开或崩溃时消费组会触发重平衡把所有分区重新分一次。重平衡期间该组的消费会停顿这也是很多线上卡顿的元凶。4.2 offset 提交机制自动提交的甜蜜陷阱消费者读消息的进度靠 offset 标记。offset 表示这个分区我已经读到哪一条了它会被提交到 Kafka 内部的__consumer_offsets主题里。新消费者加入时根据 offset 决定从哪里继续读。Kafka 消费者默认是自动提交offset 的也就是enable.auto.committrue。看着很方便但有个致命问题自动提交是周期性触发的不是处理完一条提交一条。假如你poll()拉到 500 条消息刚处理到 100 条时进程崩溃了而自动提交还没来得及执行那下次重启时这 500 条消息会全部重新消费一遍。如果处理逻辑没有幂等就会产生大量重复数据。所以可靠消费的关键是手动提交Duration timeout Duration.ofMillis(1000); while (true) { ConsumerRecordsString, String records consumer.poll(timeout); for (ConsumerRecordString, String record : records) { process(record); // 业务处理 } consumer.commitSync(); // 处理完一批再提交 offset }手动提交还有两种姿势commitSync同步提交处理完一批后阻塞等待提交成功最稳但吞吐最低commitAsync异步提交不阻塞但失败时不会自动重试需要自己监听回调。我习惯的做法是业务量小时用commitSync业务量大时用commitAsync加回调重试总之把提交和拉取之间的边界掰清楚。4.3 重复消费的根因与幂等方案重复消费不是 bug而是至少一次语义下的正常现象。Kafka 无法保证每条消息绝对只被消费一次它只能做到不丢消息但在网络抖动、消费者崩溃、提交失败时允许重复。所以后端设计必须遵循一条铁律下游处理要幂等。幂等的方法按场景选利用业务唯一键比如每条消息里带 orderId处理时先查数据库存在就跳过否则插入。数据库唯一索引是终极防线。利用 Redis setnx用消息 ID 作为 keySETNX成功才处理并设置过期时间防止并发重复进来。本地去重短时间窗口内维护已处理 ID 集合适合对内存敏感度不高的场景。我曾经接手过一个账务系统消费方一开始没做幂等对账那天突然冒出大量重复单据排查下来就是消费者在 offset 提交前崩溃重启后把上批消息重放了一遍。从那以后我对所有 Kafka 消费代码的要求都是可以重复但重复的结果必须和一次执行完全一致。4.4 多线程消费下如何保住消息顺序Kafka 的机制决定了一个分区内的消息是有顺序的但前提是只有一个线程在处理这个分区。如果你把消费到的消息丢给线程池并行处理顺序立刻被打乱。热词里那个消费端多线程如何保证消息顺序性答案就藏在这里。核心思路不是不加多线程而是让同一 key 的消息永远只交给同一个线程。做法是消费者拉完一批消息后按 key 哈希取模把消息路由到固定编号的 worker 线程int workerCount 8; ExecutorService[] workers new ExecutorService[workerCount]; for (int i 0; i workerCount; i) { workers[i] Executors.newSingleThreadExecutor(); } for (ConsumerRecordString, String record : records) { int index record.key().hashCode() % workerCount; workers[index].submit(() - process(record)); }同一个 key 哈希结果相同永远路由到同一个单线程 worker这样既利用多线程提升了吞吐又保住了同一个用户的顺序。这里要注意hashCode()的结果可能为负数取模前先Math.abs或加掩码另外至少要保证 worker 数是固定的否则改了线程数同一个 key 就被分到不同线程顺序就乱了。4.5 poll 循环里最常见的错误消费者 API 看似只有poll()一个动作但它不是随叫随到的接口而是一个持续的心跳循环。Kafka 要求消费者在max.poll.interval.ms时间内至少完成一次poll()否则会被判定为挂掉强制踢出消费组。于是很多新手会遇到这种诡异现象处理业务花了一分钟代码还在跑消息消费却是戛然而止日志里出现 rebalance。根因是默认的max.poll.interval.ms只有 5 分钟如果你的单条消息处理逻辑跑得比这还慢或者一次poll()拉了太多消息处理不完就超时了。解法有两个方向一是调大max.poll.interval.ms二是在业务侧缩短单次 poll 处理时间——比如把max.poll.records调小到 100处理完立刻commitSync再去下一轮poll()。我通常两者结合毕竟拉取和处理本来就不该在同一个循环里纠结到底。5. 从能跑到跑稳上线后的排查笔记与工具链5.1 实战排查消费组怎么一直在重复消费现象数据库里出现大量重复的订单消息消费组日志显示一直在消费但消息位点不前进。我的排查顺序是这样先看kafka-consumer-groups.sh --bootstrap-server localhost:9092 --group my-group --describe观察LAG这个指标。如果 LAG 一直在涨说明消费赶不上生产如果 LAG 为 0 但还在重放八成是提交问题。检查消费者的enable.auto.commit如果是 true再把提交间隔拉长看是否出现重复窗口把日志里每批消息的处理耗时和commit调用时间对齐能清晰看到处理完到提交前这段窗口。最后检查业务处理有没有幂等保护。没有的话先补幂等救命再思考怎么把提交窗口缩短。这个案例给我的教训是排查重复消费先别怀疑 Kafka 丢消息先怀疑自己提交姿势不对。5.2 实战排查消息延迟高到底卡在哪一段消息延迟高是使用 Kafka 最常遇到的投诉。延迟不等于 Kafka 慢而是从生产端产生消息到消费端真正处理完这段链路上某一个环节拖了后腿。我一般从三段入手生产端看生产者发送耗时和record-queue-time如果本地linger.ms调得太大或者发送失败重试次数多延迟会被强行拉高。Broker 端看 broker 机器 CPU、磁盘 IOKafka 虽然是顺序写但磁盘满了或者页缓存被其他服务挤占吞吐照样崩。消费端看消费者单条处理耗时这是延迟高最常见的位置。处理逻辑里查数据库、调远程接口都可能是瓶颈。再核对一下分区数和消费者数——如果主题 20 个分区消费组只有 2 个消费者那每个消费者平均要扛 10 个分区处理不过来LAG 自然越拉越高延迟当然大。这里有个顺手的排查技巧用kafka-consumer-groups.sh看 LAG 之后再在业务日志里给每条消息加生产时间戳和消费时间戳两边一减哪个阶段耗时最长立刻就暴露了。没有埋点延迟问题就只能靠猜。5.3 实战排查rebalance 风暴如何收场有一类故障特别迷惑人消费者没报错但整个消费组每隔几分钟就暂停一次日志里反复出现 rebalance。原因通常是某个消费者处理太慢心跳超时被踢出组里其他人接手它的分区后它又重新加入组来回拉扯这就是 rebalance 风暴。处理方案分四步调大session.timeout.ms和max.poll.interval.ms给慢消费者更多喘息空间。调小max.poll.records减小单次 poll 处理量减少超时概率。把同步的耗时业务挪出 poll 主线程改成异步处理或批量处理。给消费组接入指标监控把 rebalance 次数和 LAG 一起盯起来。但也要注意一味调大超时参数只能缓解不能根治。真正的问题是业务处理能力不够该扩容的扩容该拆分的拆分该上线程池的上线程池。5.4 集群化部署的参数底线单机版和集群之间不只是多起几台机器的问题。集群部署时至少要守住这几条底线项建议副本因子3至少 2min.insync.replicas2acks生产者设 allBroker 数至少 3奇数Controller 配置高可用机器避免和日志盘争抢Broker 数设 3 主要是为了在挂一台节点时剩下的节点还能凑够多数派完成选举和副本同步。副本因子 3 和min.insync.replicas2配合acksall意味着一条消息必须写入 2 个以上副本才算成功这样单点故障时消息不会丢。注意副本因子不能比 Broker 数还大比如 3 台机器配副本因子 5分区永远无法同步日志里全是Not enough replicas报错。集群安装本身不复杂每台机器下载相同版本的 Kafkaserver.properties里配置broker.id必选唯一KRaft 模式下用同一个cluster id格式化存储目录再逐个启动就行。难的是启动顺序和配置一致性建议用配置管理工具统一维护别手敲。5.5 可视化工具和命令行三板斧很多初学者拿到 Kafka 不知道从哪看消息这里分享我常用的三件套命令行工具Kafka 自带kafka-topics.sh、kafka-console-consumer.sh、kafka-consumer-groups.sh。我每天用得最多的是kafka-consumer-groups.sh --describe一眼看尽每个消费组的 LAG。Offset Explorer原 Kafka Tool桌面客户端适合快速浏览主题、查看分区和消息内容。注意它只适合开发环境别在压测环境乱连。Kafka UI开源项目Web 面板比桌面工具好在团队共享可以直接在浏览器里查主题、看消息、查看消费者组 LAG。后端把 Kafka 节点暴露到内网团队成员就都可以自助排查。至于消息内容是否可视化我有个提醒生产环境的消息里常有敏感数据用户 ID、手机号给开发环境开可视化没问题生产环境最好只给运维和核心开发开并且做好权限控制。5.6 顺带一提C/Qt 项目中怎么接 Kafka如果你是在 C 或 Qt 项目里对接 Kafka绕不开的是 librdkafka。它是 C/C 生态的事实标准客户端Confluent 官方也在维护。核心配置和 Java 端一致RdKafka::Conf *conf RdKafka::Conf::create(RdKafka::Conf::CONF_GLOBAL); conf-set(bootstrap.servers, localhost:9092, errstr); conf-set(group.id, my-group, errstr); RdKafka::KafkaConsumer *consumer RdKafka::KafkaConsumer::create(conf, errstr);Windows 上如果要用 MinGW 编译 librdkafka最容易踩的坑是动态库版本必须和你编译器位数一致x64 的 Qt 程序不能链 x86 的 librdkafka另外 librdkafka 依赖 openssl 和 zlib要提前装好并让 CMake 能找到它们。Qt 本身不提供 Kafka 组件所以整体思路就是Qt 写界面和业务librdkafka 负责和 Kafka 通信两边用信号槽接起来。代码上尽量把 librdkafka 的消费回调放到独立线程避免卡住 Qt 的事件循环。6. 面试必问这五个 Kafka 原理题最好别答错6.1 为什么 Kafka 快得像不像传统消息队列Kafka 的高吞吐不是靠复杂的缓存算法而是靠几个非常朴素的底层设计顺序写磁盘传统消息队列表面上是存内存最终落地时其实是随机写盘。Kafka 直接以追加形式顺序写文件磁盘顺序写比随机写快几个量级。页缓存Kafka 不自己做缓存而是依赖操作系统的页缓存。刚刚写入的数据消费者往往马上能读到命中的是内存不是磁盘。零拷贝消费端读数据时Kafka 通过sendfile把磁盘数据直接传给网卡省掉拷到用户态的环节。批量与压缩生产者攒一批再发消费者一批一批拉网络传输成本被均摊。我把这套设计理解为快递集散中心散户一件一件发货物流成本高、效率低Kafka 把大量包裹先按目的地打包再用大车统一运输时间没有少太多但吞吐量完全不是一个量级。6.2 ISR、副本和 acks可靠性的三个齿轮副本不是越多越好关键是 Kafka 怎么判断一个 follower 是不是跟得上。Kafka 用 ISRIn-Sync Replica集合来管理ISR 里的副本必须持续从 leader 同步数据差距超过阈值就会被踢出去。这里有三层配置联动producer 的acks决定要等几个副本确认broker 的min.insync.replicas决定 ISR 里至少要有几个副本副本因子replication.factor决定每个分区有几个副本。打个比方一个分区有 leader 和 2 个 followerISR3min.insync.replicas2。假设一个 follower 扛不住被踢出 ISRISR 变成 2 个此时写请求依然成功因为 2 个节点满足最小同步副本要求如果另一个节点也挂了ISR 只剩 1 个再写就会报NotEnoughReplicasException宁可拒绝写入也不丢消息。6.3 Leader 选举到底怎么选一个分区的多个副本分为 leader 和 follower读写全走 leader。那 leader 挂了选谁答案是在 ISR 里选。为什么强调 ISR因为 ISR 里的副本已经和 leader 保持了同步数据最全选它不会丢消息。如果 ISR 里的副本都挂了Kafka 会从虽然不在 ISR 但还在线的副本中挑一个这时可能出现数据丢失是无奈之下的恢复手段。节点间的协调由 Controller 负责集群中会选出一个 Controller Broker专门处理分区分配、副本变动这些元数据操作。6.4 消费组重平衡的完整过程重平衡的本质是消费组成员变了分区重新分配。完整过程大致是消费者启动或退出时向组协调者发送加入组请求协调者从组里挑一个消费者当 leader把成员列表和订阅信息给它leader 负责制定分配方案把分区分给每个成员协调者把方案下发给所有成员各自开始领取对应分区。重平衡期间所有成员都会停止消费所以它越频繁系统吞吐就越差。面试时如果能补充重平衡可能由消费者处理超时触发的吗并给出解决方案会比单纯背流程加分。6.5 从生产者到消费者消息不丢失的完整链路如何保证 Kafka 消息不丢失是面试高频题答的时候必须覆盖整条链路生产者到 Brokeracksall让消息必须写进多个副本才算成功同时开启重试和幂等发送失败自动重发且不产生重复。Broker 内部副本因子设 3min.insync.replicas2容忍单机故障刷盘策略保持默认Kafka 靠副本同步保证持久性而不是靠 fsync。Broker 到消费者消费者关掉自动提交处理完再手动提交 offset处理动作要做幂等防止重复消息带来脏数据。三兄弟缺一不可生产者不丢是源头副本不丢是存储层消费者不丢是末端兜底。少一个环节不丢失就是一句空话。我自己带新人时最后总会说一句Kafka 的 API 其实两小时就能学会难的是把消息丢失、重复消费、顺序错乱、延迟升高这些概念内化成系统设计时的直觉。先搭一个单机版写一段生产者、一段消费者故意杀掉消费者进程看重复消费怎么发生再用命令行工具盯一次 LAG 的变化把这些场景亲手跑过一遍比背十篇面试题都管用。