Kafka 面试题
Kafka 面试题
导语:Kafka 的面试重点在高吞吐的底层原理、副本与 ISR 机制、HW/LEO、Rebalance 与位移管理、Exactly-Once,以及 KRaft 的版本演进。其中「
acks=all必须配合min.insync.replicas」和「Rebalance 的触发与避免」是最容易被问倒的两处。共 16 题。
一、架构与高性能
1. Kafka 的整体架构由哪些组件组成?
答:
| 组件 | 职责 |
|---|---|
| Producer | 生产者,负责发送消息;可指定 key 来决定落到哪个分区 |
| Broker | Kafka 服务节点,负责消息的存储与读写 |
| Topic | 消息的逻辑分类 |
| Partition | Topic 的分区,是并行与扩容的基本单位;同一分区内消息有序 |
| Replica | 分区的副本,分 Leader(提供读写)与 Follower(只负责同步) |
| Consumer Group | 消费组,组间广播、组内竞争 |
| Controller | 集群中被选举出来的一个 Broker,负责分区 Leader 选举、副本状态管理、元数据变更分发 |
| Controller Quorum / ZooKeeper | 元数据与选主的协调者(Kafka 4.0 起只有 KRaft,不再支持 ZK) |
三条核心约束(必记):
- 一个分区在同一消费组内只能被一个消费者消费 → 消费者并发上限 = 分区数;
- 同一分区的多个副本必须分布在不同的 Broker 上,否则副本失去意义;
- 顺序性只在分区级别成立——跨分区无序。
2. Kafka 为什么吞吐量最高?
答: 六个设计共同作用,缺一不可:
| 设计 | 原理 | 收益 |
|---|---|---|
| ① 顺序写磁盘 | 日志 append-only,磁盘顺序写可达数百 MB/s,接近随机内存访问的速度 | 把"磁盘慢"这个前提消掉 |
| ② 页缓存(Page Cache) | Kafka 自己不管理缓存,直接复用 OS 的 Page Cache:写入即入页缓存并返回,由 OS 异步刷盘 | 避免 JVM 堆内存的 GC 压力;读写都命中缓存时几乎无磁盘 IO |
| ③ 零拷贝(sendfile) | 消费时用 FileChannel.transferTo(底层 sendfile),数据从 Page Cache 直接 DMA 到网卡 | 不经用户态、不做 CPU 拷贝,省掉 2 次 CPU 拷贝与 2 次模式切换 |
| ④ 批量 + 压缩 | Producer 累积成批发送(batch.size、linger.ms),Broker 批量落盘;支持 gzip / snappy / lz4 / zstd | 用"延迟几十毫秒"换"吞吐数量级提升" |
| ⑤ 分区并行 | 一个 Topic 分多个 Partition,分布在不同 Broker 上 | 读写都在分区级别并行,可水平扩展 |
| ⑥ 稀疏索引 | 为日志段建立稀疏索引(.index / .timeindex) | 索引文件小(可常驻内存),定位快 |
纠正一处常见表述:不要写"Kafka 把消息写入堆外内存"。Kafka 用的是操作系统的 Page Cache(由内核管理,与 JVM 无关),这也是它比"JVM 内缓存"方案更抗 GC 的原因。Kafka 自身只在零拷贝路径上涉及堆外(由 JNI 直接传递)。
3. Kafka 的物理存储结构是怎样的?
答: 三级结构:Topic → Partition → Log Segment。
目录结构(每个分区是一个目录,命名 <topic>-<partition>):
/kafka-logs/order-topic-0/
├── 00000000000000000000.log # 消息数据文件(默认 1GB 一段)
├── 00000000000000000000.index # 偏移索引(offset → 物理位置)
├── 00000000000000000000.timeindex # 时间戳索引(timestamp → offset)
├── 00000000000000000487.snapshot # 幂等生产者/事务相关的快照
├── leader-epoch-checkpoint # Leader Epoch 记录(防止 HW 截断问题)
└── partition.metadata # 分区元数据日志段(Log Segment)的滚动与清理:
- 当
.log达到log.segment.bytes(默认 1GB)或超过log.roll.hours(默认 168h)时,切出新段; - 旧段按保留策略整段删除(
log.retention.hours默认 168h /log.retention.bytes)——"整段删除"是 Kafka 清理数据极其高效的原因(只删文件,不用逐条删)。
稀疏索引的设计要点(面试重点):
- 索引不是每条消息都建,而是每写入
log.index.interval.bytes(默认 4KB)才记录一条索引; - 查找流程:先用二分法在
.index中定位到"≤ 目标 offset 的最近索引项",拿到物理位置后,从该位置开始顺序扫描.log直到找到目标 offset; - 为什么用稀疏索引:若为每条消息建索引,索引文件会与数据文件一样大,无法全量加载到内存;稀疏索引让索引文件足够小(可常驻 Page Cache),用"一次小范围顺序扫描"换取"索引常驻内存"。
追问:"按时间查消息怎么实现?"——用
.timeindex找到"该时间点对应的 offset",再按 offset 定位。所以按时间查询的本质是先转成 offset 查询。
日志段为什么设计成固定大小(1GB):便于用文件偏移做 mmap 与快速定位,也让"删除过期数据"变成"删除文件"这一 O(1) 操作。
二、可靠性:副本与不丢消息
4. Kafka 如何保证消息不丢失?
答: 三层设防,但有一个必须同时配置才有意义的组合(这是本题的核心考点):
① 生产端
| 配置 | 作用 |
|---|---|
acks=all(-1) | 要求 ISR 中所有副本都写入才返回成功 |
retries > 0 | 允许重试(配合幂等生产者避免重试造成重复/乱序) |
enable.idempotence=true | 幂等生产者(Kafka 3.0+ 默认开启),通过 PID + 序列号去重 |
② Broker 端
| 配置 | 作用 |
|---|---|
replication.factor >= 3 | 至少 3 个副本 |
min.insync.replicas >= 2 | 关键!见下方说明 |
unclean.leader.election.enable=false | 禁止非 ISR 副本当选 Leader(默认就是 false),防止已提交消息被回退 |
③ 消费端
enable.auto.commit=false,先处理业务、成功后再手动提交位移;- 拿到消息就自动提交位移(或先提交再处理),一旦处理失败就会丢消息。
⚠️ 最重要的组合关系:acks=all 必须配合 min.insync.replicas >= 2
这是面试中最容易答错的地方:
acks=all的语义是"ISR 中所有副本都写入才成功";- 但如果 ISR 里只剩 Leader 一个副本(其他副本都因滞后被踢出 ISR),那么"ISR 的所有副本"就只剩 Leader——此时
acks=all退化成acks=1!Leader 写入即返回,Leader 随后宕机,消息就丢了; min.insync.replicas=2的作用:当 ISR 中可用副本数少于该值时,生产者写入会直接失败(抛NotEnoughReplicasException),而不是"假装成功"。
这就是"宁可写失败,也不要静默丢数据"的设计——先保证不会误以为成功,再由生产端重试或告警。
一句话总结:
acks=all保证"写进 ISR 全部副本",min.insync.replicas保证"ISR 不会退化到只剩 1 个副本",两者缺一不可。
5. Kafka 中 ack=0/1/-1 的区别?
答:
| acks | 语义 | 可靠性 | 吞吐 |
|---|---|---|---|
0 | 生产者发出即认为成功,完全不等待 Broker 响应 | 最低:网络丢包、Broker 宕机都无感知 | 最高 |
1 | Leader 写入本地日志即返回成功 | 中:Leader 在副本同步之前宕机就丢数据 | 高 |
-1(即 all) | ISR 中所有副本都写入才返回 | 最高 | 最低(延迟也最高) |
三点补充:
- 从可靠性排序:
all > 1 > 0;从吞吐排序:0 > 1 > all——这就是"可靠性与吞吐"的经典权衡; - Kafka 3.0 起,生产者的默认
acks从1改为all(配合默认开启的幂等生产者),官方倾向于"默认更安全"; acks=all单独配置是不够的,必须配合min.insync.replicas(见上一题)。
6. 什么是 ISR、OSR、AR?
答:
| 概念 | 含义 |
|---|---|
| AR(Assigned Replicas) | 分区的所有副本集合 = ISR + OSR |
| ISR(In-Sync Replicas) | 与 Leader 保持同步的副本集合,包含 Leader 自身 |
| OSR(Out-of-Sync Replicas) | 被踢出 ISR 的滞后副本 |
ISR 的判定标准(容易答错):
- 判定依据是 Follower 是否在
replica.lag.time.max.ms(默认 30s)内持续向 Leader 发起 fetch 请求; - ⚠️ 不是"落后多少条消息"(旧版本曾用
replica.lag.max.messages,早已废弃); - 只要 Follower 持续拉取(哪怕暂时落后但没超过时间阈值),就仍在 ISR 中;追上进度后可以重新加入 ISR(ISR 是动态的)。
ISR 的三大作用:
acks=all的"all"指的是 ISR,不是 AR——所以可以通过"缩小 ISR"来降低写入延迟(代价是可靠性下降);- Leader 选举只在 ISR 中进行(
unclean.leader.election.enable=false时),从而保证"已提交消息不丢"; - 它定义了"已提交(committed)"的含义——只有被 ISR 全部同步的消息才算提交。
参数调优的权衡:
replica.lag.time.max.ms调大 → 慢副本留在 ISR(写成功更快但可靠性下降);调小 → 副本频繁进出 ISR,增加 Leader 切换风险。这是一个典型的"延迟 vs 可靠"权衡点。
7. 什么是 HW(高水位)与 LEO?
答:
| 概念 | 定义 | 记忆 |
|---|---|---|
| LEO(Log End Offset) | 每个副本下一条将要写入消息的位移,即该副本日志的末尾位置 | "我写到哪了" |
| HW(High Watermark) | 所有 ISR 副本都已同步到的位置,即 min(所有 ISR 副本的 LEO) | "大家确认到哪了" |
三条关键规则:
- 消费者只能读到 HW 之前的消息——HW 之后的消息虽然已写入 Leader,但尚未被所有 ISR 副本确认,因此对消费者不可见;
- HW 定义了"已提交(committed)"消息的边界——"已提交"= "已复制到所有 ISR 副本";
- Leader 切换时以 HW 为参考判断哪些消息是"确定不会丢"的。
为什么需要 HW:保证"消费者读到的消息一定不会因为 Leader 故障切换而消失"。若消费者能读到 HW 之后的数据,而该 Leader 随后宕机、新 Leader 恰好缺少这部分数据,就会出现"读到过的数据又没了"的不一致。
HW 的隐患与 Leader Epoch 的引入(加分点):
- Kafka 早期存在著名的"HW 截断导致数据不一致/丢数据"问题:Leader 切走后,由于各副本 HW 更新存在延迟,副本在重新成为 Leader 时依据自己记录的 HW 做日志截断,可能截掉了本该保留的数据,导致同一 offset 在不同副本上的数据不一致;
- 解决方案:引入 Leader Epoch(记录在
leader-epoch-checkpoint文件中)——每个 Leader 任期对应一个 Epoch,新 Leader 通过 Epoch 与起始 offset 准确判断"应该保留到哪个位置、从哪里开始截断",不再依赖 HW 做截断决策,从而消除了这类不一致。
一句话记忆:LEO 是"本地写入位置",HW 是"ISR 公认位置",消费者只能看 HW 之前;而"截断"这件事交给 Leader Epoch 更可靠。
三、顺序、位移与 Exactly-Once
8. Kafka 如何保证消息顺序?
答: 先说结论:Kafka 的顺序性只在分区内成立,跨分区不保证。
标准做法(两步):
- 生产端:同一业务 key 路由到同一分区
- 发送时指定
key(如orderId),Kafka 用hash(key) % partitionCount计算分区(默认分区器); - 这样同一订单的消息必然进入同一分区,分区内天然有序;
- 若 key 为 null,则按"粘性分区"策略轮流分配,不保证有序。
- 发送时指定
- 消费端:单线程消费该分区
- 一个分区在本消费组内只会分配给一个消费者,该消费者对它串行消费即保证顺序;
- 如果用多线程处理,必须按 key 分桶(同一 key 交给同一线程),否则会乱序。
两个会破坏顺序的坑(高频追问):
| 坑 | 原因 | 对策 |
|---|---|---|
| 生产者重试导致乱序 | max.in.flight.requests.per.connection > 1 且开启重试时,"请求 1 失败重试、请求 2 先成功",分区内顺序就乱了 | 开启幂等生产者 enable.idempotence=true(Kafka 3.0+ 默认开启)——Broker 会按序列号排序并去重,保证分区内严格有序且不重复;或把 in-flight 设为 1(性能损失大) |
| 增加分区数导致乱序 | hash(key) % N 中的 N 变了,同一 key 的新消息会落到不同分区,与历史消息不在同一分区 | 规划分区数时预留余量,避免频繁扩分区 |
9. Kafka 如何实现 Exactly-Once(精确一次)语义?
答: 两个组件组合实现:
① 幂等生产者(Idempotent Producer)
- 开启
enable.idempotence=true(Kafka 3.0+ 默认开启); - 原理:Producer 初始化时从 Broker 获取一个 PID(Producer ID),每条消息带上
(PID, 分区, 序列号); - Broker 为每个
(PID, Partition)维护"已接收的最大序列号",对重复序列号直接丢弃,从而避免"重试导致的单分区重复"; - 局限:只能保证"单分区 + 单生产者会话内"的去重。
② 事务(Transactions)
- 开启
transactional.id(必需,且建议每个生产者实例唯一),通过initTransactions()获取 PID; - API:
beginTransaction()/commitTransaction()/abortTransaction(); - 引入 Transaction Coordinator(由某个 Broker 担任),事务状态写入内部 Topic
__transaction_state; - 事务消息会写入"控制消息(Control Batch)"标记该事务的提交/回滚;
- 消费者侧:设置
isolation.level=read_committed,则只读取已提交事务的消息(未提交事务的消息会被跳过);read_uncommitted(默认)则能看到未提交的消息。
组合起来实现 EOS 的经典场景:consume-transform-produce(从 Kafka 读 → 处理 → 写回 Kafka),把"消费位移提交"也纳入同一个事务,从而实现"要么这批数据被处理并写出、位移也提交;要么全部回滚重来"。
⚠️ 最重要的澄清(高频追问):
Kafka 的 EOS 是"Kafka 内部"的精确一次。它保证的是"Kafka → 处理 → Kafka"这条链路不重不丢。端到端是否精确一次,还取决于下游系统:
- 若下游是外部 DB/HTTP 接口 → 仍需业务幂等;
- 因为"写外部系统"和"提交 Kafka 位移"无法组成一个原子事务——这与第《消息队列通用》中"端到端恰好一次靠幂等实现"的结论一致。
10. 什么是 Rebalance?触发条件与危害是什么?
答: Rebalance 指消费组内分区所有权在消费者之间重新分配的过程,由 Group Coordinator(组协调者,由某个 Broker 担任)负责协调。
三种触发条件:
| 触发条件 | 说明 |
|---|---|
| ① 消费者数量变化 | 新消费者加入(joinGroup)或消费者主动关闭(close) |
| ② 订阅的 Topic 分区数变化 | 例如给 Topic 增加了分区 |
| ③ 消费者心跳超时 / 处理超时 | 消费者未在 session.timeout.ms(默认 45s)内发送心跳,被认为"已死亡";或两次 poll() 的间隔超过 max.poll.interval.ms(默认 5min)——后者是生产上最常见的 Rebalance 原因(单批消息处理太慢) |
四大危害:
- Stop The World:Rebalance 期间整个消费组暂停消费,等待新分配完成;分区多、消费者多时可能持续数十秒到分钟级;
- 消费滞后(Lag 飙升):停顿期间消息持续堆积,可能触发下游延迟告警甚至事故;
- 重复消费:某个分区的位移未及时提交就被转交给新消费者,新消费者从旧位移开始消费 → 必须靠幂等兜底;
- 恶性循环:Rebalance 拖慢消费 →
poll间隔更长 → 更容易超时被踢出 → 再次 Rebalance。
如何避免或减少 Rebalance:
| 手段 | 说明 |
|---|---|
调大 max.poll.interval.ms,或减小 max.poll.records | 确保"一批消息能在间隔内处理完"——最有效的组合 |
合理设置 session.timeout.ms | 太小会因网络抖动误判死亡;太大会导致故障发现慢 |
使用静态成员 group.instance.id | 消费者重启时不触发 Rebalance(滚动发布场景收益极大) |
使用 CooperativeStickyAssignor(协作式再均衡) | 只迁移需要变动的分区,无需"全体撤销再重新分配",大幅缩短停顿 |
| 消费者实例数不要超过分区数 | 多出来的实例会空闲,且其进出会触发 Rebalance |
避免长时间阻塞 poll 线程 | 把耗时处理交给线程池,poll 线程只负责拉取 |
11. Kafka 有哪些分区分配策略?
答:
| 策略 | 机制 | 特点 |
|---|---|---|
| Range(范围) | 按 Topic 维度:把每个 Topic 的所有分区按范围均分给消费者 | 简单;但当"分区数不能被消费者数整除"时,靠前的消费者会多分一个分区,多个 Topic 叠加后倾斜明显 |
| RoundRobin(轮询) | 把所有订阅 Topic 的全部分区排序后轮流分配 | 比 Range 更均匀;但消费者订阅的 Topic 不一致时仍可能不均 |
| Sticky(粘性) | 在尽量保留原有分配的前提下做均衡 | Rebalance 时迁移的分区最少,减少停顿与重复消费 |
| CooperativeSticky(协作式粘性,官方推荐) | 在 Sticky 基础上采用增量式再均衡 | 只迁移需要变动的分区,避免 Eager 模式的"全员放弃再重分配",Rebalance 影响最小 |
两个必须区分的概念(高频):
| 模式 | 行为 | 影响 |
|---|---|---|
| Eager(积极/重平衡协议 v1) | Rebalance 时所有消费者先放弃全部分区,再重新分配 | 全组 STOP THE WORLD,停顿长 |
| Cooperative(协作式/协议 v2) | 只撤销/迁移需要变动的分区,其他分区继续消费 | 停顿极短,是新版本推荐 |
配置方式:partition.assignment.strategy(消费者端配置),可传多个策略做协商,如 CooperativeStickyAssignor。
实践建议:新集群直接用
CooperativeStickyAssignor;从老版本升级时注意它是"消费者协议"的变更,需要按滚动升级流程操作(否则可能触发多次 Rebalance)。
12. 消费者位移(offset)是如何提交与存储的?
答:
① 提交方式(两种)
| 方式 | 配置与 API | 特点 |
|---|---|---|
| 自动提交 | enable.auto.commit=true(默认值)+ auto.commit.interval.ms(默认 5s) | 由 poll() 在后台自动提交上一次 poll 返回的位移。⚠️ "提交时机"与"业务处理完成"无关 → 可能在业务处理前就提交了位移,崩溃即丢消息 |
| 手动提交(推荐) | enable.auto.commit=false,业务处理完成后调用 | ① commitSync():同步提交,阻塞但会重试,可靠性高;② commitAsync():异步提交,不阻塞但失败不重试(可能丢位移);常见组合:处理中异步提交、关闭前用同步提交兜底 |
② 存储位置:__consumer_offsets 内部 Topic
- 位移不是存在 ZK(老版本才存 ZK,早已废弃),而是存在 Kafka 内部的
__consumer_offsetsTopic 中(默认 50 个分区,offsets.topic.replication.factor默认 3); - 存储格式:以
(group.id, topic, partition)为 Key,offset 为 Value 写入;由 Group Coordinator 负责管理; - 关键理解:"提交位移"本质上就是"向
__consumer_offsets生产一条消息"——所以它会经过正常的 Kafka 写入路径,有延迟、也可能失败。这也是为什么"自动提交"不能完全信任、"同步提交"更可靠的原因。
③ 位移重置:auto.offset.reset(最容易误解的参数)
| 取值 | 行为 |
|---|---|
earliest | 从最早的位移开始消费(会重复消费历史数据) |
latest | 从最新开始(默认值,会跳过历史数据) |
none | 找不到位移时直接抛异常 |
⚠️ 关键澄清:
auto.offset.reset只在"该消费组没有已提交位移"或"位移越界(消息已被删除)"时生效——它不是"每次启动都重置"。若消费组已有提交过的位移,这个参数完全不起作用。想真正"重置已有消费组的位移",必须用
kafka-consumer-groups.sh --reset-offsets(且要求该消费组处于非活跃状态),或换一个新的group.id。
四、高可用与 KRaft
13. Kafka 的消息保留与日志压缩策略?
答: 由 cleanup.policy 控制,两种策略解决两个不同的问题:
| 策略 | 机制 | 适用场景 |
|---|---|---|
delete(默认) | 按时间(log.retention.hours,默认 168h = 7 天)或大小(log.retention.bytes)整段删除过期的日志段 | 普通消息:日志、埋点、事件流 |
compact(日志压缩) | 对同一个 key,只保留最新的 value,更早的记录在压缩时被清理 | 变更日志 / 状态快照:__consumer_offsets、Kafka Streams 的状态存储、CDC 全量快照同步 |
| 组合 | cleanup.policy=compact,delete | 既按 key 保留最新,又按时间清理 |
日志压缩(Log Compaction)的机制与要点:
- 后台 Log Cleaner 定期扫描,把"同一 key 的旧版本"标记为可清理,异步执行(不是实时,因此某个 key 可能短时间内存在多个版本);
- 保留 tombstone(墓碑消息):value 为
null的消息是"删除标记",会额外保留一段时间(delete.retention.ms,默认 24h)后才真正清理——这样下游才能"读到删除事件"; - 正在被消费的活跃段不会被压缩(
min.cleanable.dirty.ratio控制压缩触发比例); - 压缩后 offset 不连续(被删记录的 offset 会出现空洞)——但 offset 本身不会变化,这是"位移可回溯"的前提。
⚠️ 高频混淆点:"日志压缩(Log Compaction)" ≠ "压缩算法(compression)"
- Log Compaction:按 key 去重,保留最新值,目的是"让 topic 成为一份 key 的最新状态快照";
- Compression(gzip/snappy/lz4/zstd):字节级压缩,目的是减少网络与磁盘占用。
这两个概念在面试中经常被混为一谈,务必明确区分。
14. Controller 的作用是什么?
答: Controller 是集群中被选举出来的一个 Broker,承担集群级的"管理与协调"职责(相当于集群的"大脑"):
| 职责 | 说明 |
|---|---|
| ① 分区 Leader 选举 | 当某个分区的 Leader 所在 Broker 宕机,Controller 从该分区的 ISR 中选出新 Leader,并把新的 Leader/Follower 信息下发给相关 Broker |
| ② 副本状态管理 | 维护分区的 ISR、HW 等元数据 |
| ③ Topic 元数据变更分发 | Topic 的创建/删除、分区数增加等操作的编排与分发 |
| ④ Broker 上下线处理 | Broker 加入/退出时的分区重新分配与通知 |
选举机制(两种模式):
| 模式 | 选举方式 |
|---|---|
| ZooKeeper 模式(旧) | 每个 Broker 争抢在 ZK 上创建 /controller 临时节点,创建成功者成为 Controller;Controller 宕机后临时节点消失,其他 Broker 重新竞选 |
| KRaft 模式(新) | 由 Controller Quorum(一组专门的 Controller 节点,基于 Raft)选举产生,元数据以事件日志形式在 Quorum 中复制 |
两个关键认知:
- Controller 自身无状态(元数据都在 ZK/KRaft 中),因此它宕机不会丢数据,新 Controller 秒级选出并从 ZK/KRaft 重建内存状态;
- 但"重建元数据"的成本随分区数增长——分区数达到数十万时,Controller 故障恢复可能耗时数分钟,期间分区 Leader 无法切换(故障恢复能力瘫痪)。这正是 KRaft 要解决的核心痛点之一。
15. Kafka 的高可用是如何实现的?
答: 一套组合拳,五个要素:
| 要素 | 配置/机制 | 作用 |
|---|---|---|
| ① 分区多副本 | replication.factor = 3 | 每个分区有 3 个副本,分布在不同 Broker(同 Broker 上的副本无意义) |
| ② Leader/Follower 模型 | 读写都走 Leader;Follower 主动 fetch 拉取同步 | Follower 只做备份,不承担读写(保证一致性) |
| ③ ISR 机制 | replica.lag.time.max.ms 判定 | 只有"同步中"的副本参与 acks=all 与 Leader 选举 |
| ④ Controller 自动选主 | Leader 宕机 → Controller 从 ISR 选新 Leader | 秒级完成切换,无需人工介入(这是相对 RocketMQ 传统主从的优势) |
| ⑤ 防止数据回退 | min.insync.replicas=2 + unclean.leader.election.enable=false | 前者保证"副本不足时写失败",后者禁止"落后副本当选 Leader" |
生产推荐配置("容忍 1 节点故障且不丢数据"):
replication.factor=3 # 3 副本
min.insync.replicas=2 # ISR 至少 2 个副本才允许写入
acks=all # 生产端等待 ISR 全部确认
unclean.leader.election.enable=false # 禁止非 ISR 副本当选为什么是这套组合:3 副本 + min.insync.replicas=2 意味着最多容忍 1 个 Broker 宕机——宕机后 ISR 还剩 2 个副本,仍满足"≥2",写入继续可用且不丢数据;若再挂一台,ISR 只剩 1 个,写入直接报错(保护数据而不是静默降级)。
两个补充点:
replication.factor只能增加、不能减少(需用kafka-reassign-partitions.sh做副本重分配);- 副本越多写放大越严重(每份数据要写 N 次),所以副本数不是越多越好,3 是业界共识的平衡点。
16. Kafka 为什么要去 ZooKeeper 化(KRaft)?
答: ZooKeeper 模式的四大痛点:
| 痛点 | 说明 |
|---|---|
| ① 元数据吞吐受限(最核心) | ZK 是 CP 系统,元数据变更需多数派确认。分区数达到数十万时,Controller 故障后的元数据重建可能耗时数分钟,期间集群"失去故障切换能力" |
| ② 羊群效应(Herd Effect) | 每个 Broker 在 ZK 上注册大量 Watcher,一次状态变更会唤醒大量监听器,ZK 容易过载甚至成为瓶颈 |
| ③ 架构冗余、状态双写 | 元数据既存在 ZK,又在 Kafka Controller 内存中维护一份,两套状态需保持一致,增加了复杂度与不一致风险 |
| ④ 运维成本 | 需要额外维护一套独立的 ZK 集群(ZK 自身也要 3/5 节点做高可用),部署与故障排查都更复杂 |
KRaft 的解决思路:
- 用 Kafka 内置的 Raft 实现(Kafka Raft Metadata Quorum) 自管元数据;
- 元数据以事件日志的形式存储,并通过 Raft 复制到一组专门的 Controller 节点(Controller Quorum);
- Controller 选举由 Raft 完成,不再需要 ZK;
- 支持元数据快照,新节点可快速追平,而不用回放全部日志。
收益:
| 收益 | 说明 |
|---|---|
| 支持百万级分区 | 元数据吞吐不再受 ZK 限制 |
| Controller 故障恢复秒级 | 元数据在 Raft 日志中,恢复无需重建全量内存状态 |
| 集群启动更快 | 不再依赖 ZK 的连接与 Watcher 注册 |
| 少维护一套系统 | 运维复杂度显著下降 |
版本演进(必须记准,高频考点):
| 版本 | 状态 |
|---|---|
| Kafka 2.8 | KRaft 作为预览特性引入 |
| Kafka 3.3 | KRaft 生产可用(GA) |
| Kafka 3.5 / 3.6 | ZooKeeper 模式被标记弃用,官方推动迁移 |
| Kafka 4.0 | 完全移除 ZooKeeper 支持,KRaft 成为唯一模式 |
面试提醒:不要再说"Kafka 依赖 ZooKeeper"。准确表述是:"Kafka 早期依赖 ZooKeeper 做元数据与选主;2.8 引入 KRaft,3.3 生产可用,3.5+ 弃用 ZK,4.0 已彻底移除 ZK。"这个版本演进是能直接体现技术敏感度的加分点。
