Skip to content

Kafka 面试高频题总索引

本文件汇总 21 章节面试题,按主题分类提炼(不重复每章原文),共 11 大类、110 + 道,全部带「考察点 / 标准答案 / 加分项 / 易错点」四件套。

💻 使用建议

  1. 先看下方「目录脑图」,定位你想冲的方向。
  2. 每类至少把前 5 题刷透,然后挑剩下题里带 ⭐⭐⭐⭐ 的高频考点啃下来。
  3. 文末「高频题 Top 30」直接对标大厂一面。

🎯 难度图例

  • ⭐ 入门必会(基础概念 / 命令)
  • ⭐⭐ 进阶常考(原理理解 / 选型)
  • ⭐⭐⭐ 大厂高频(生产细节 / Kafka vs 其他 MQ)
  • ⭐⭐⭐⭐ 资深岗位(深入内核 / 故障排查 / EOS / KRaft)

🗺 目录脑图


一、架构与基础概念

关联章节:Ch 01 / Ch 02 / Ch 03。 这一类是几乎所有岗位的「入场题」——答不上来直接挂。

1.1 Topic / Partition / Offset / Broker / Consumer Group 分别是什么?三句话讲清楚关系 ⭐

  • 考察点:Kafka 的 5 个核心概念是否真的搞懂。
  • 标准答案
    1. Topic:消息的逻辑分类(如「订单事件」),生产者写到 Topic、消费者读 Topic。
    2. Partition:Topic 内部的物理切片,是顺序保证 / 并行度 / 副本 / 故障转移的最小单元;同一个 Topic 由 N 个 Partition 组成。
    3. Offset:单个 Partition 内的单调递增 64 位序号,每条消息有唯一 offset,Consumer 用 offset 标记读到哪。
    4. Broker:一个 Kafka 服务进程,1 ~ N 个 Broker 组成集群;每个 Partition 的副本分布在不同 Broker。
    5. Consumer Group:一组 Consumer 共享一个 group.id;组内每个 Partition 只被一个成员消费,组间互不影响(一份数据多吃)。
  • 加分项:补一句「Partition 是分布式的原子单位——Kafka 几乎所有能力(顺序、副本、Rebalance、ISR、EOS、Compaction)都围绕 Partition 展开」。
  • 易错点:把 Topic 当成 RabbitMQ 的 Queue(Topic 是逻辑、Partition 才是物理);Offset 是 partition 级而不是 topic 级。

1.2 Kafka 是「队列」还是「日志」?为什么这么说有意义? ⭐⭐

  • 考察点:是否理解 Kafka 的设计哲学。
  • 标准答案
    1. Kafka 的本质是分布式提交日志(Distributed Commit Log),而不是「消息队列」。
    2. 队列 = 出队即消失;日志 = 追加写 + 任意位置可重放。
    3. 因此 Kafka 同时支持「MQ 异步 / 削峰 / 解耦」+「Event Sourcing / CDC / 实时数仓 / 历史回放」多种场景。
  • 加分项:举例「同一份订单事件被风控、库存、通知、数仓四个 Consumer Group 各读一份;新业务上线还能从头读半年的历史」——这是队列模型做不到的。
  • 易错点:把 Kafka 当成「带持久化的 RabbitMQ」。

1.3 一条消息从 Producer 到 Consumer 的完整链路? ⭐⭐⭐

  • 考察点:能否串起整条数据流。
  • 标准答案(按时序):
    1. Producer 应用 send → 序列化 → Partitioner 选 Partition → 入 Producer 端 RecordAccumulator 攒 batch
    2. Sender 线程批量发往 Partition Leader 所在 Broker(NetworkClient → SocketServer)
    3. Leader 写本地日志(append 到 active segment)→ 等待 Follower 拉副本
    4. Follower 通过 FetchRequest 拉数据,Leader 收到 Follower 推进 → 更新 HW
    5. 根据 acks 给 Producer 返回 ack
    6. Consumer poll → FetchRequest 到 Leader → Leader 返回 ≤ HW 的数据 → 反序列化 → 业务处理 → commit offset 到 __consumer_offsets
  • 加分项:补「整条链路里 Broker 始终是被动响应;Producer 攒批 + Consumer 拉模式是 Kafka 高吞吐的关键」。
  • 易错点:忘记 HW 这道「可见性」门槛;以为 Broker 主动推送给 Consumer。

1.4 一个分区可以同时被同一消费组内的多个 Consumer 消费吗?为什么? ⭐⭐

  • 考察点:消费组的「分区独占」语义。
  • 标准答案
    1. 不可以。在同一个 Consumer Group 内,一个 Partition 只能被一个 Consumer 实例消费
    2. 这是为了保证单 Partition 内顺序消费和 offset 提交的一致性。
    3. 不同 Consumer Group 之间互不影响(广播消费)。
  • 加分项:「这也意味着 Consumer 数量超过 Partition 数时,多余的 Consumer 是闲置的——所以 Partition 数 = 消费并行度上限」。
  • 易错点:把分区独占跟「整个 Topic 只能被一个 Consumer 消费」搞混。

1.5 Kafka 集群里 Coordinator 和 Controller 有什么区别? ⭐⭐⭐

  • 考察点:两个「协调者」角色容易混。
  • 标准答案
    • Group Coordinator:每个 Consumer Group 一个,负责 Rebalance / 心跳 / offset 提交,本质是某个 Broker 上的一段逻辑(哪个 Broker 上托管了 group 对应的 __consumer_offsets 分区,谁就是 Coordinator)。
    • Controller:整个集群一个(KRaft 时代是 Quorum 选 Active Controller),管 Topic / Partition / 副本元数据、Leader 选举、ISR 变更。
  • 加分项:「Coordinator 是『消费侧』的协调者;Controller 是『集群层』的协调者,两者从不重叠」。
  • 易错点:以为 Controller 也管 Consumer Group。

1.6 Kafka 客户端是怎么知道集群拓扑的?Broker 挂了客户端会怎样? ⭐⭐⭐

  • 考察点:Metadata 协议 + 客户端容错。
  • 标准答案
    1. 客户端启动时连 bootstrap.servers任意一个 Broker → 发 MetadataRequest → 拿全集群 Topic / Partition / Leader 元数据。
    2. 客户端本地缓存元数据,按 metadata.max.age.ms(默认 5min)刷新,或在收到 NotLeaderForPartitionException 时主动刷。
    3. Broker 挂了 → 客户端发请求收到错误 → 触发刷新元数据 → 拿到新 Leader → 重试。
  • 加分项:bootstrap.servers 写多个的意义是「初始化容错」,不需要写全部 Broker。
  • 易错点:以为 Producer 每次 send 都查元数据。

1.7 Kafka 为什么不像 RabbitMQ 那样「服务端推消息」? ⭐⭐⭐

  • 考察点:拉 vs 推的取舍。
  • 标准答案
    1. 拉模型让消费速率由消费者决定,避免快生产者打挂慢消费者。
    2. 拉模型支持批量请求 + 长轮询,吞吐效率高。
    3. 拉模型让消费位点(offset)由客户端管理,方便重放 / 回溯
    4. 推模型在多消费方时控速难(要做反压协议),Kafka 设计初衷是高吞吐数据总线。
  • 加分项:补「不过纯拉的代价是空轮询,Kafka 用 fetch.min.bytes + fetch.max.wait.ms 做长轮询缓解」。
  • 易错点:以为拉模型只有缺点。

1.8 ZK 和 KRaft 时代的元数据存储有什么区别? ⭐⭐⭐

  • 考察点:架构演进。
  • 标准答案
    • ZK 时代(≤ 2.x):Topic / 分区 / Broker / ACL 等全在 ZooKeeper;Controller 选举、元数据广播都依赖 ZK;扩展性瓶颈在 ZK Watch + 单 Controller 推送。
    • KRaft 时代(2.8 Preview / 3.3 GA / 4.0 完全移除 ZK):元数据自己写到 __cluster_metadata 内部 Topic,由 Controller Quorum(Raft)维护;事件溯源 + 快照恢复;摆脱 ZK 后单集群可支持百万分区Broker 启动加速
  • 加分项:补 KIP-500 设计目标 + Migration KIP-866 ZK→KRaft 的迁移路径。
  • 易错点:把 KRaft 当成 Kafka 自己塞了一个 ZK——KRaft 是 Raft 实现,根本不依赖第三方协调服务。

1.9 Kafka 集群常见的「内部 Topic」有哪些? ⭐⭐

  • 考察点:是否做过运维。
  • 标准答案
    • __consumer_offsets:消费组提交的 offset,cleanup.policy=compact,副本因子建议 3
    • __transaction_state:事务元数据,cleanup.policy=compact
    • __cluster_metadata:KRaft 元数据日志(替代 ZK)
    • _schemas:Schema Registry 存放 schema(compact)
    • connect-configs / offsets / status:Kafka Connect 元数据(compact)
  • 加分项:「这些 Topic 副本数应当 ≥ 3,绝对不要手动删——删了消费组 Lag 全丢、事务永久卡死」。
  • 易错点:把 _schemas 算成「业务 Topic」。
  • 考察点:定位与生态位。
  • 标准答案
    • Kafka 是数据流的总线(pipe):日志聚合用它做 buffer(Filebeat → Kafka → Logstash → ES);流处理用它做输入输出(Kafka → Flink → Kafka);数仓用它做实时入仓通道(MySQL → Debezium → Kafka → ClickHouse / Doris)。
  • 加分项:「Kafka 在 Linkedin 诞生就是为了取代之前『N 系统对 N 系统点对点同步』的混乱,把所有数据汇总成总线」。
  • 易错点:把 Kafka 跟 Flink 当成竞品(其实是搭档)。

二、存储与日志格式

关联章节:Ch 07 存储 / Ch 08 高吞吐底层 / Ch 14 Compaction

2.1 Topic-Partition-Segment 三层结构是什么?为什么要分段? ⭐⭐⭐

  • 考察点:存储布局。
  • 标准答案
    1. Topic → N 个 Partition(物理切片)→ 每个 Partition 是一个目录,包含多个 Segment(日志段文件)。
    2. 每个 Segment 由一组同名文件组成:<baseOffset>.log(数据)+ .index(offset 索引)+ .timeindex(时间索引)+ 可选 .snapshot / .txnindex
    3. 分段的目的:老数据按文件批量删除/压缩(按时间或大小过期整段就行,不用扫单条);索引文件 mmap 加速查找写入只针对 active segment,rollover 时更换。
  • 加分项:滚 segment 触发条件 log.segment.bytes(默认 1GB)/ log.roll.ms / log.roll.hours
  • 易错点:以为一个 Partition 就一个文件。

2.2 .index 是稠密索引还是稀疏索引?怎么找一条 offset 对应的消息? ⭐⭐⭐⭐

  • 考察点:索引数据结构。
  • 标准答案
    1. 稀疏索引.index 不为每条记录建索引,每隔 index.interval.bytes(默认 4KB)建一条「offset → 文件物理位置 position」的映射。
    2. 查找 offset 时:① 二分定位 segment 文件(按 baseOffset);② 在 .index 上二分找到不大于目标 offset 的最大索引项;③ 从该 position 起顺序扫 .log 直到目标 offset。
    3. .timeindex 同理:「时间戳 → offset」稀疏映射,配合 .index 完成「按时间找消息」。
  • 加分项:稀疏索引 + mmap = 索引文件常驻 PageCache,查找几乎零 I/O。
  • 易错点:以为是 B-Tree 或哈希索引。

2.3 V0 / V1 / V2 RecordBatch 格式有什么差别?为什么 V2 是分水岭? ⭐⭐⭐

  • 考察点:协议演进。
  • 标准答案
    • V0 / V1(0.10 之前):每条消息独立 → 每条都有 CRC、Magic、Attributes、key/value 长度字段,头部开销大;V1 加了 timestamp。
    • V2(0.11 起):引入 RecordBatch(一批消息共享元信息:baseOffset / partitionLeaderEpoch / producerId / producerEpoch / baseSequence ...),单条 record 只存增量 → 头部开销大幅降低;天然支持 Header(k=v 元数据)幂等 / 事务(PID + epoch + seq 字段就在 batch 头里)。
  • 加分项:V2 让端到端压缩率更好(共享 batch 头)+ CRC 在 batch 级而非 record 级,少做一次 CRC。
  • 易错点:以为 timestamp 一直存在;以为事务跟消息格式无关。

2.4 顺序写为什么比随机写快几个数量级? ⭐⭐⭐

  • 考察点:底层 I/O 原理。
  • 标准答案
    1. 机械盘随机写需要寻道 + 旋转 + 写 → 几 ms/次;顺序写只移动磁头连续吐数据 → 几 μs/次,相差 1000 倍。
    2. SSD 上随机写虽快,但仍有写放大、垃圾回收问题;顺序写让 FTL 友好。
    3. Kafka 把 partition 当顺序日志写,加上 PageCache 批量回写,写放大 ≈ 1
  • 加分项:「LinkedIn 的早期论文测过,6 块 SATA 盘做 RAID0 顺序写能到 600MB/s,比单块随机写快两个数量级」。
  • 易错点:以为「内存随机比磁盘顺序快」——错,对于大文件顺序追加,磁盘顺序吞吐能匹敌内存随机。

2.5 PageCache 在 Kafka 里扮演什么角色?为什么 Broker 不在 JVM 堆里缓存消息? ⭐⭐⭐⭐

  • 考察点:操作系统 + JVM 设计权衡。
  • 标准答案
    1. 写入:Producer 数据 → Broker write() → 进 PageCache → OS 异步刷盘。
    2. 读取:Consumer Fetch → 命中 PageCache 直接返回(绝大多数 offset 是「最近写入」,热数据 100% 命中)。
    3. 不在 JVM 堆缓存的原因:① 避免重复缓存(OS PageCache 已经缓存了一次);② 避免堆变大触发 Full GC;③ Broker 进程崩溃时数据仍在 PageCache,OS 完成刷盘,重启后还能恢复。
  • 加分项:「这就是 Kafka 推荐 6~10G JVM 堆 + 操作系统留 50%+ 内存给 PageCache 的根因」。
  • 易错点:以为 Kafka 自己实现了 LRU 缓存。

2.6 零拷贝 sendfile 在 Kafka 中怎么用?省掉了什么? ⭐⭐⭐⭐

  • 考察点:Linux 系统调用 + 性能优化。
  • 标准答案
    1. 传统 read+write 路径:磁盘 → 内核 PageCache → 用户态缓冲 → 内核 Socket buffer → 网卡,4 次拷贝 + 2 次用户/内核态切换
    2. sendfile / transferTo 路径:磁盘 → 内核 PageCache → 网卡(或 DMA gather copy),2 次拷贝 + 1 次切换;用户态完全不出现数据。
    3. Kafka 在 Consumer Fetch 响应时用 FileChannel.transferTo()(最终调用 sendfile),把 segment 文件直接发给客户端 socket。
  • 加分项:「这是 Kafka 能把单 Broker 出口流量推到 GB/s 的关键,否则用户态 buffer 会成为瓶颈」。
  • 易错点:以为「零拷贝 = 完全没有拷贝」(其实仍有 2 次内核内拷贝)。

2.7 mmap 在 Kafka 里用在哪里?为什么不用于 .log 文件? ⭐⭐⭐⭐

  • 考察点:mmap 的适用场景。
  • 标准答案
    • mmap 用在 .index / .timeindex 索引文件:索引小、查找频繁,mmap 后整个文件映射到虚拟内存,二分查找 = 内存访问。
    • 不用于 .log:① log 文件大(GB 级),mmap 会污染地址空间;② sendfile 已经解决了 log 文件读取性能;③ mmap 写入需要应用做 msync,控制刷盘麻烦。
  • 加分项:补「Kafka 还用 mmap 做 __cluster_metadata 快照恢复」。
  • 易错点:以为整个 Kafka 都靠 mmap。

2.8 Log Compaction 的工作机制是什么?什么场景必须用? ⭐⭐⭐⭐

  • 考察点:Compaction 内核 + 应用场景。
  • 标准答案
    1. 机制:后台 LogCleaner 线程把同一 Key 的旧版本压缩掉,只保留最新值null value = Tombstone(墓碑),表示删除。
    2. 触发min.cleanable.dirty.ratio(默认 0.5)+ segment.ms,仅 active 之外的 segment 才会被 compact。
    3. 保留语义:active segment 不动;非 active 的 segment 中同 key 多条 → 合并为一条最新值。
    4. 必用场景:状态快照(K-V 表)、CDC 变更流(最新行状态)、__consumer_offsets / __transaction_state 内部 Topic。
  • 加分项:补「cleanup.policy=compact,delete 组合:compact 同时按 retention.ms 整段删旧 segment,适合『最近 30 天的 K-V 状态』」。
  • 易错点:以为 Compaction 会立即生效;以为 active segment 也参与。

2.9 deletecompact 两种 cleanup.policy 何时选什么? ⭐⭐⭐

  • 考察点:保留策略选型。
  • 标准答案
    • delete(默认):按时间 / 字节滚段删除整个 segment,适合事件流 / 日志 / 时序数据
    • compact:按 Key 保留最新值,适合状态快照 / 内部元数据 Topic
    • compact,delete:兼具两者,适合「保留最近 30 天的状态快照」。
  • 加分项:举 __consumer_offsets 必须 compact 的例子(否则旧 group 的 offset 永远不丢)。
  • 易错点:把 compact 当成「定时清空」。

2.10 一条消息从落盘到对消费者可见,要经历哪些「门槛」? ⭐⭐⭐

  • 考察点:HW 可见性边界。
  • 标准答案
    1. Leader 写本地 → LEO 推进到 N+1
    2. Follower 拉到这条消息 → 它们的 LEO 也推进
    3. Leader 收到所有 ISR 的 Follower fetch → 计算 min(LEO) → 推进 HW
    4. Consumer Fetch 时 Broker 只返回 offset < HW 的消息
  • 加分项:「所以未达到 HW 的消息对 Consumer 是『不可见』的——这是 Kafka 保证『副本一致性』的基础」。
  • 易错点:以为 Broker 落盘就可见。

三、生产者

关联章节:Ch 04 Producer 深入 / Ch 13 EOS

3.1 acks=0/1/all 三种语义有什么区别?分别在什么时候用? ⭐⭐⭐

  • 考察点:可靠性级别。
  • 标准答案
    • acks=0:不等任何 ack,发完即返回,吞吐最高 / 易丢;适合指标埋点 / 日志
    • acks=1:Leader 写本地即返回;中等可靠;Leader 切换瞬间可能丢;适合容忍少量丢失的业务
    • acks=all:等 ISR 内所有副本确认(结合 min.insync.replicas ≥ 2);最可靠;适合关键业务(订单 / 支付 / 风控)
  • 加分项:「3.0+ 默认 acks=all;但 min.insync.replicas 默认 1 会让 acks=all 形同虚设,必须显式调到 2」。
  • 易错点:以为 acks=all 就一定不丢。

3.2 幂等 Producer 是怎么实现的?解决了什么问题? ⭐⭐⭐⭐

  • 考察点:幂等机制内核。
  • 标准答案
    1. 解决:网络重试导致同一条消息被 Broker 写两次的问题。
    2. 实现:开 enable.idempotence=true 后,Broker 给 Producer 分配一个 PID(Producer ID)+ 每条消息按 partition 维度递增的 sequence number;Broker 端按 (PID, seq) 去重。
    3. 限制:单分区有效;Producer 进程重启后 PID 变化,幂等失效(这就要事务来兜底)。
  • 加分项:3.0+ 默认 enable.idempotence=true;它会自动设 acks=all + max.in.flight ≤ 5 + retries=Integer.MAX
  • 易错点:以为幂等能防业务侧重复(业务侧重复要业务幂等去重)。

3.3 linger.msbatch.size 的关系?怎么调? ⭐⭐⭐

  • 考察点:批次机制。
  • 标准答案
    • batch.size(默认 16KB):单分区批次的字节上限;攒满即发。
    • linger.ms(默认 0):批次未满时最多等多久;时间到也发。
    • 真实触发条件:任一满足即发
  • 加分项:「吞吐优先:linger.ms=20~50 + batch.size=128KB~256KB;延迟优先:linger.ms=0 + batch.size=16KB」;监控 batch-size-avg 接近上限说明攒得满。
  • 易错点:单独调 batch.size 不调 linger.ms 没意义;把 batch.size 当 byte 总量(其实是「单分区」的)。

3.4 Producer 分区策略有哪些?默认是什么? ⭐⭐⭐

  • 考察点:路由机制 + 演进。
  • 标准答案
    • 有 Key:murmur2(key) % numPartitions同 Key 落同分区保证顺序。
    • 无 Key:① 老版本(< 2.4)= 轮询(每条切换分区,batch 难攒满);② 2.4 ~ 现在 = Sticky Partitioner(粘住一个分区直到 batch 满 / linger 到,再换),显著提升吞吐。
    • 自定义:实现 org.apache.kafka.clients.producer.Partitioner 接口。
  • 加分项:3.3+ 引入 partitioner.ignore.keysBuiltInPartitioner 改进,避免低流量时 Sticky 把流量集中。
  • 易错点:以为默认是「随机分区」。

3.5 5 种压缩算法(none / gzip / snappy / lz4 / zstd)怎么选? ⭐⭐⭐

  • 考察点:压缩选型。

  • 标准答案

    算法CPU压缩率推荐场景
    none0内网巨大带宽
    gzip2.5~3.5×离线归档
    snappy1.5~2×兼容老系统
    lz41.5~2×延迟优先(推荐)
    zstd2~3×吞吐+存储双优先(强推)
  • 加分项:压缩在 Producer 端做,Broker 不解压(直接存压缩后的 batch),Consumer 端解压 → 端到端节省网络 + 磁盘

  • 易错点:选 gzip 在高 QPS 下打满 CPU。

3.6 Producer 怎么保证「Exactly Once」? ⭐⭐⭐⭐

  • 考察点:EOS 链条。
  • 标准答案
    1. 第一层:幂等 Producer → 防同一 Producer 进程网络重试导致的重复。
    2. 第二层:事务 Producer → 防 Producer 进程重启后 PID 变化、跨分区原子写。
    3. 第三层:消费者 isolation.level=read_committed → 过滤未提交事务的消息。
    4. 第四层:消费者业务幂等 → 兜底(消费 N 次结果一致)。
  • 加分项:「Kafka 的 EOS 严格意义只对『consume-process-produce』链路有效(Streams / Connect 模式);对『外部系统副作用』(写 DB / 发 HTTP)只能通过业务幂等保证」。
  • 易错点:以为开了事务就万事大吉。

3.7 transactional.id 怎么取?为什么不能跨实例共用? ⭐⭐⭐⭐

  • 考察点:事务身份与 PID Fencing。
  • 标准答案
    1. transactional.id 是事务的全局唯一身份;必须跨重启稳定,重启后用同 ID 触发 PID Fencing 把旧 Producer 挤掉。
    2. 跨实例共用会互相 Fence:A 启动后 B 启动 → A 报 ProducerFencedException 退出;A 又重启 → B 退出 → 死循环。
    3. 推荐:tx-<service>-<pod-name> / tx-<service>-<host>-<port>,K8s StatefulSet 拿稳定 podName。
  • 加分项transaction.timeout.ms 必须 ≤ Broker 的 transaction.max.timeout.ms(默认 15min)。
  • 易错点:把 transactional.id 写成 "my-tx" 静态字符串。

3.8 Producer 端 delivery.timeout.ms / request.timeout.ms / retries 三者是什么关系? ⭐⭐⭐

  • 考察点:超时模型。
  • 标准答案
    • request.timeout.ms(默认 30s):单次请求超时。
    • retries(默认 Int.MAX):失败重试次数。
    • delivery.timeout.ms(默认 2min):从 send() 到成功 / 失败的总时长上限这是一个硬上限,超过即使 retries 还没用完也会失败。
  • 加分项:「大批量 Producer 推荐把 delivery.timeout.ms 调到 5min,避免突发慢盘 / 慢网导致大批 send 失败」。
  • 易错点:只看 retries,忘了 delivery.timeout.ms。

3.9 Producer 如何处理 Topic 不存在 / Leader 切换的暂时错误? ⭐⭐

  • 考察点:错误分类与重试。
  • 标准答案
    1. 可重试错误LeaderNotAvailableException / NotLeaderForPartitionException / NetworkException / RetriableException 子类 → 客户端自动重试。
    2. 不可重试错误InvalidTopicException / RecordTooLargeException / OffsetOutOfRangeException → 直接抛给业务。
    3. 重试时按 retry.backoff.ms(默认 100ms)退避。
  • 加分项:错误码到异常的映射在 org.apache.kafka.common.errors 包下都有。
  • 易错点:以为所有错误都自动重试。

3.10 一个 Producer 实例可以并发地往多个 Partition 发消息吗?是怎么做的? ⭐⭐

  • 考察点:客户端线程模型。
  • 标准答案
    1. Producer 是线程安全的,业务多线程可共享同一个 KafkaProducer 实例。
    2. 内部 send() 是异步:消息进 RecordAccumulator(按分区分桶),后台 Sender 线程批量发往各分区 Leader。
    3. 不同 Broker 用不同 socket;同一 Broker 上的多分区可共享 socket(NetworkClient)。
  • 加分项:「这就是为什么官方 best practice 是『每应用一个 KafkaProducer 实例』,而不是每线程一个」。
  • 易错点:以为每发一条都阻塞。

四、消费者与消费组

关联章节:Ch 05 Consumer 深入 / Ch 11 Rebalance / Ch 12 消息语义

4.1 subscribeassign 有什么区别? ⭐⭐

  • 考察点:自动 vs 手动管理分区。
  • 标准答案
    • subscribe(topics)加入消费组,由 Coordinator 通过 Rebalance 协议分配分区;offset 提交到 __consumer_offsets
    • assign(partitions)手动指定消费哪些分区,不参与 Rebalance;不能与 subscribe 混用。
  • 加分项:assign 适合「单实例消费 + 自定义 offset 存储」(如 Spark Streaming Direct API 早期)。
  • 易错点:以为 assign 也会自动管理 offset。

4.2 poll() 循环模型是怎么工作的? ⭐⭐⭐

  • 考察点:Consumer 主循环。
  • 标准答案
    1. 第一次 poll:发现没有 group 元数据 → JoinGroup → SyncGroup → 拿到分配的分区 → 拉位点
    2. 每次 poll:① 心跳(自动后台线程发);② Fetch 数据;③ 触发 OffsetCommit(自动提交时);④ 返回 ConsumerRecords
    3. 业务处理 records → 下次 poll 之前必须完成(否则触发 max.poll.interval.ms
  • 加分项:「心跳走独立线程(0.10.1 起),但 max.poll.interval.ms 要在主线程 poll 调用之间衡量——否则就算心跳还在,业务卡住也会被踢」。
  • 易错点:以为心跳跟 poll 是一回事。

4.3 自动提交 vs 手动提交,什么时候选什么?怎么避免坑? ⭐⭐⭐⭐

  • 考察点:offset 提交时机。
  • 标准答案
    • 自动提交(enable.auto.commit=true):定时(auto.commit.interval.ms)后台提交,与业务无关 → 业务挂了 offset 已提交 → 丢消息。
    • 手动提交(commitSync / commitAsync):业务处理完成后再提交。
    • 推荐关掉自动提交 + 业务处理完手动提交
  • 加分项:commitSync 阻塞、保证;commitAsync 非阻塞、可能漏;常见模式是「正常用 async,关闭前用 sync 兜底」。
  • 易错点:用自动提交 + 「为了防止丢手动提一遍」——其实只要开了自动提交,定时器还会在后台跑。

4.4 Rebalance 的触发条件有哪些? ⭐⭐⭐

  • 考察点:协议事件。
  • 标准答案
    1. 消费组成员变化:新成员加入 / 心跳超时 / 主动 LeaveGroup
    2. 订阅 Topic 变化:subscribe 列表变 / 通配匹配新增 Topic
    3. Topic 元数据变化:分区数增加(不能减少)
    4. Coordinator 变更:托管 group 的 __consumer_offsets 分区 Leader 切换
  • 加分项:补「Static Membership 能让短暂重启不触发 Rebalance;Cooperative 让 Rebalance 不再 STW」。
  • 易错点:以为 Rebalance 只在 Consumer 起停时发生。

4.5 4 种分配策略(Range / RoundRobin / Sticky / CooperativeSticky)有什么区别? ⭐⭐⭐⭐

  • 考察点:分配策略选型。
  • 标准答案
    • Range(默认之一):按 Topic 单独算,分区数 / 消费者数 + 余数给前面 → 同一 Topic 内分区数不均时前面消费者多分
    • RoundRobin:跨 Topic 轮询分配 → 每个消费者拿到的分区数最多差 1。
    • Sticky(0.11+):尽量让上次分配保持不变,新一轮在最少改动的前提下重新平衡。
    • CooperativeSticky(2.4+):Sticky 的增量协作版本,Rebalance 时不放弃所有分区(只移交部分),避免 Stop-the-World。
  • 加分项:3.0+ 默认 [Range, CooperativeSticky];推荐统一切到 CooperativeSticky
  • 易错点:以为切策略可以一步到位(要按 KIP-429 走双策略迁移)。

4.6 Eager Rebalance 与 Cooperative Rebalance 的本质差别? ⭐⭐⭐⭐

  • 考察点:Rebalance 协议演进。
  • 标准答案
    • Eager(老):Rebalance 触发 → 所有 Consumer revoke 全部分区 → JoinGroup → SyncGroup → 拿新分配 → 重新开始消费。期间所有人Stop-the-World
    • Cooperative(KIP-429):分两阶段;第一阶段只 revoke 「不再属于自己的分区」,其余继续消费;第二阶段把 revoke 的分区转给新 owner。最小化停顿
  • 加分项:「这是 Kafka Streams、Connect 默认采用的协议」;KIP-848 新一代 Consumer Group Protocol 进一步把协议从客户端搬到 Broker 侧。
  • 易错点:以为 Cooperative 完全没有停顿。

4.7 Static Membership 是什么?解决了什么问题? ⭐⭐⭐

  • 考察点:KIP-345。
  • 标准答案
    1. 给 Consumer 设 group.instance.id=<稳定唯一> → Coordinator 把这个 id 当成员身份。
    2. 短暂重启(< session.timeout.ms)时,Coordinator 不踢,新进程接管旧成员的分区,不触发 Rebalance
  • 加分项:K8s StatefulSet 拿稳定 podName 即可;session.timeout.ms 调到 60s 给重启留时间。
  • 易错点:以为 Static Membership = 永远不 Rebalance(实际仍会在真长时间下线时触发)。

4.8 Consumer Lag 怎么定义?怎么排查持续增长? ⭐⭐⭐

  • 考察点:监控 + 排障。
  • 标准答案
    • Lag = LogEndOffset - 已提交 offset,按 partition 维度计算。
    • 排查 6 步:① 业务处理耗时(CPU / 下游 RT);② 消费者实例数 vs 分区数(实例少了);③ 是否在 Rebalance(last-rebalance-seconds-ago);④ Broker Side 是否吞吐瓶颈(fetch 慢);⑤ 反序列化 / 网络异常(records-consumed-rate 是否为 0);⑥ Topic 流量是否突增。
  • 加分项:监控 records-lag-max 客户端 JMX 比 Broker 端 kafka-consumer-groups --describe 实时性更好。
  • 易错点:把 Lag 当「时间」而非「条数」。

4.9 seek / pause / resume 的典型用法? ⭐⭐

  • 考察点:高级 API。
  • 标准答案
    • seek(partition, offset):跳到指定 offset(重放历史 / 跳过坏消息)。
    • seekToBeginning / seekToEnd:跳到头 / 尾。
    • offsetsForTimes(map<tp,timestamp>):按时间戳找 offset → 配 seek 实现「按时间消费」。
    • pause(partitions) / resume(partitions):暂停 / 恢复某些分区拉取(不影响心跳),用于背压。
  • 加分项pause 的典型场景是「下游限流 / DB 卡 → 暂停消费让 Lag 增长,业务恢复后 resume」。
  • 易错点:以为 pause 会触发 Rebalance(不会,自己不放分区)。

4.10 同一个消费组内不同实例订阅了不一致的 Topic 列表,会发生什么? ⭐⭐⭐

  • 考察点:Group 协议约束。
  • 标准答案
    • Coordinator 会取所有实例订阅 Topic 的并集,但不一致是反模式:分配策略对不同 Topic 计算结果不一致,可能导致部分实例闲置。
    • 同时不同实例的 partition.assignment.strategy 不一致 → 直接 InconsistentGroupProtocolException
  • 加分项:「实战中所有 Consumer 实例必须用同一份配置 + 同一份订阅列表」。
  • 易错点:以为可以「分组消费一组 Topic」(请用不同 group.id)。

五、副本与 ISR

关联章节:Ch 09 副本与 ISR

5.1 ISR 是什么?谁可以进 / 出 ISR? ⭐⭐⭐

  • 考察点:ISR 定义。
  • 标准答案
    1. ISR (In-Sync Replicas) = 与 Leader 保持同步的副本集合(含 Leader 自身)。
    2. 进入条件:Follower 追上 Leader(LEO 与 Leader LEO 之差 ≤ 0 且最近一次 fetch 时间 < replica.lag.time.max.ms)。
    3. 出 ISR:超过 replica.lag.time.max.ms(默认 30s)没追上。
  • 加分项:「老版本(0.9 之前)还有 replica.lag.max.messages 按消息数判定,但很容易在突发流量下抖动,0.9 起删除」。
  • 易错点:以为 ISR 包括所有副本。

5.2 HW 和 LEO 是什么?为什么需要 HW? ⭐⭐⭐⭐

  • 考察点:可见性边界。
  • 标准答案
    • LEO (Log End Offset) = 副本下一条要写入的 offset(每个副本各自的 LEO)。
    • HW (High Watermark) = ISR 中所有副本 LEO 的最小值只有 < HW 的消息对 Consumer 可见
    • HW 的作用:防止 Leader 切换后「Consumer 看到的消息丢失」——因为新 Leader 的 LEO 至少 ≥ 旧 HW(旧 HW 之前的数据所有 ISR 都有)。
  • 加分项:HW 只在 Leader 上权威推进,Follower 通过 fetch 响应感知(旧版本通过单独 propagate)。
  • 易错点:HW 不是「写盘水位」,是「可见水位」。

5.3 Leader Epoch 解决了什么问题? ⭐⭐⭐⭐

  • 考察点:副本一致性深坑。
  • 标准答案
    1. 问题:旧版本 Kafka 在 Leader 切换 + Follower 截断到 HW 时存在「数据不一致 / 数据丢失」窗口。
    2. 机制:每次 Leader 切换 epoch 加 1;每个 Follower 维护 (epoch → start offset) 映射;Follower 截断时用 OffsetsForLeaderEpoch 请求新 Leader 拿到「该 epoch 的最大 offset」,按这个截断而非 HW。
    3. 效果:避免「截到 HW 但旧 Leader 实际还有未达成 ISR 共识的数据被错误保留」的隐患。
  • 加分项:KIP-101 / KIP-279。
  • 易错点:把 Leader Epoch 跟 Controller Epoch 搞混。

5.4 unclean.leader.election.enable 开了会怎样? ⭐⭐⭐⭐

  • 考察点:可用性 vs 可靠性。
  • 标准答案
    1. :ISR 全挂时允许非 ISR 副本当 Leader → 该副本数据落后 → 已 ack 的消息「凭空消失」。
    2. (2.0+ 默认):ISR 全挂时分区 OFFLINE,保数据不丢。
  • 加分项:「数据库 / 金融 / 订单类业务必须关;一些低价值数据(埋点 / 日志聚合)可考虑开来换可用性」。
  • 易错点:以为关了会让集群挂;其实只是少数分区不可写。

5.5 min.insync.replicas + acks=all 的组合保证了什么? ⭐⭐⭐⭐

  • 考察点:可靠性配置组合。
  • 标准答案
    1. acks=all → Producer 等 ISR 内所有副本确认。
    2. min.insync.replicas=2(3 副本时)→ ISR 不足 2 时 Broker 拒绝写(NotEnoughReplicasException)。
    3. 组合保证:至少 2 副本同时确认才认为写成功;挂 1 个副本仍可写;挂 2 个副本拒绝写(保不丢)。
  • 加分项:补「单独配 acks=allmin.insync.replicas=1 → ISR 内只有 Leader 时仍能写 → 等于 acks=1」。
  • 易错点:以为 acks=all 一招鲜。

5.6 副本同步是 Leader 推还是 Follower 拉?为什么? ⭐⭐⭐

  • 考察点:副本协议设计。
  • 标准答案
    • Follower 拉(FetchRequest)。原因:① 拉模式让 Follower 控速避免被打挂;② 复用 Consumer fetch 路径,简化协议;③ Leader 不需要维护「每个 Follower 进度」的推送状态。
  • 加分项:Follower fetch 可以选择最近副本 fetching(KIP-392,2.4+),跨机房客户端能就近拉,节省带宽。
  • 易错点:以为 Leader 主动推。

5.7 ISR 抖动(频繁收缩 / 扩张)的常见原因? ⭐⭐⭐

  • 考察点:故障排查。
  • 标准答案
    1. 网络抖动:跨机房 / 网卡限速 / TCP 重传。
    2. 磁盘抖动:fsync 慢、磁盘 IO 拥塞。
    3. GC 长停顿:Follower JVM Full GC 几秒。
    4. 配置激进replica.lag.time.max.ms 设得太小(< 10s)。
    5. Topic 流量突增:副本来不及追。
  • 加分项:监控 IsrShrinksPerSec / IsrExpandsPerSec,每分钟超过 5 次告警。
  • 易错点:直接调大 replica.lag.time.max.ms 治表不治本。

5.8 副本数 3 时,挂 1 个 / 2 个 Broker 分别会发生什么? ⭐⭐⭐

  • 考察点:故障建模。
  • 标准答案
    • 挂 1 个:① 该 Broker 上是 Leader 的分区 → Controller 选 ISR 中其他副本为新 Leader(瞬间);② ISR 大小变 2,仍 ≥ min.insync.replicas=2,业务感知短抖动。
    • 挂 2 个:① 大部分分区 Leader 切到剩下的 1 个 Broker;② ISR 只剩 1,min.insync.replicas=2 时业务写直接被拒NotEnoughReplicasException),保不丢;③ 若 unclean.leader.election=true 又有非 ISR 副本能起,可能导致丢数据。
  • 加分项:「3 副本是平衡点;金融业务可上 5 副本(容忍 2 节点同时故障)」。
  • 易错点:以为「挂 2 个仍能正常写」。

5.9 replica.lag.time.max.ms 该怎么设? ⭐⭐

  • 考察点:调参经验。
  • 标准答案
    1. 默认 30s,适合大多数场景。
    2. 抖动严重 → 调到 60s,但同时排查根因。
    3. 设得太小(< 10s)→ ISR 频繁抖动。
    4. 设得太大(> 5min)→ 落后副本长期占 ISR 名额,acks=all 等慢副本,写延迟变大。
  • 加分项:跨机房部署需要明显放大。
  • 易错点:以为调小能让数据更可靠。

5.10 一个分区的副本怎么分配到不同 Broker?rack-aware 是什么? ⭐⭐⭐

  • 考察点:副本放置策略。
  • 标准答案
    1. 默认放置:尽量打散到不同 Broker(避免单 Broker 多副本);分区 0 的 Leader 选起始 Broker,副本顺序往后。
    2. rack-awarebroker.rack 配 + 0.10+):把不同副本尽量分到不同 rack(机柜 / 可用区),同 rack 故障时仍有副本存活。
  • 加分项:「云上常用 broker.rack=us-east-1a 这种 AZ 标签 → 副本天然跨 AZ 分布」。
  • 易错点:以为副本是随机放置。

六、Controller 与 KRaft

关联章节:Ch 10 Controller 与 KRaft

6.1 ZK 时代 Controller 的职责是什么?为什么是单点? ⭐⭐⭐

  • 考察点:ZK 架构。
  • 标准答案
    1. 职责:Topic / Partition 元数据管理、Leader 选举、ISR 变更广播、Reassignment 执行。
    2. 选举:所有 Broker 抢注 ZK 临时节点 /controller,先到先得。
    3. 单点:任意时刻集群只有 1 个 Active Controller;不是单机故障点(挂了 ZK 重新选举),但变更串行化,规模上去后扩展性差。
  • 加分项:「Controller 重启时要从 ZK 全量拉元数据并广播给所有 Broker,百万分区规模下要几分钟」。
  • 易错点:以为 Controller 是个独立进程(其实是 Broker 的一个角色)。

6.2 KRaft 是什么?解决了什么 ZK 时代的痛点? ⭐⭐⭐⭐

  • 考察点:KIP-500 设计。
  • 标准答案
    1. KRaft = Kafka Raft Metadata mode:用 Raft 协议把元数据写到内部 Topic __cluster_metadata,Controller Quorum(3 / 5 节点)维护。
    2. 痛点:① 摆脱 ZK 运维负担;② 元数据事件溯源(替代 ZK 全量拉)→ Broker 启动 / Controller 切换都从增量日志 + 快照恢复,秒级完成;③ 单集群分区数从 ~20 万扩到百万。
  • 加分项:「Active Controller 只有 1 个,但 Quorum 容错 (N-1)/2;元数据日志写入也是 Raft 多数派提交」。
  • 易错点:以为 KRaft 是「Kafka 自己嵌了一个 ZK」。

6.3 KRaft 的 Raft 选举怎么工作? ⭐⭐⭐⭐

  • 考察点:Raft 协议。
  • 标准答案
    1. Controller Quorum 节点维护 term;初始或 Active 失联时进入 Candidate 状态,发 VoteRequest,过半票成为 Active Controller。
    2. Active Controller 把元数据变更(Topic 创建 / 副本分配 / ISR 变更)作为 Record 写到 __cluster_metadata Topic(Raft Replicated Log)。
    3. Broker 作为 Observer 订阅这个 Topic,按 Record 同步本地元数据状态机。
  • 加分项:补 KIP-595(Raft 协议)/ KIP-630(快照)。
  • 易错点:把 Active Controller 当成「老 Controller 角色」——它管元数据 Quorum,不再像 ZK 时代那样需要全广播。

6.4 ZK → KRaft 怎么迁移?要注意什么? ⭐⭐⭐⭐

  • 考察点:KIP-866。
  • 标准答案
    1. 先起独立的 KRaft Controller Quorum(3 节点)。
    2. Broker 配 zookeeper.metadata.migration.enable=true,仍连 ZK + 同时注册到 KRaft Controller → 此时元数据从 ZK 双写到 KRaft。
    3. 等 KRaft 完全同步 → 把 Broker 配置切到 KRaft only(去掉 ZK 配置)→ ZK 下线。
    4. 期间不要乱改 process.roles、不要清 KRaft 数据卷。
  • 加分项:「Kafka 3.5 完善了迁移工具;3.9 是最后一个支持 ZK 的版本系列;4.0 完全移除 ZK」。
  • 易错点:以为「重启 Broker 自动切换」。

6.5 KRaft 模式下 Controller Quorum 用 3 节点还是 5 节点? ⭐⭐⭐

  • 考察点:高可用规划。
  • 标准答案
    • 3 节点:容忍 1 故障,足够中小集群(< 10 Broker)。
    • 5 节点:容忍 2 故障,适合大集群 / 跨机房部署。
    • 不要 1 / 2 / 4 / 6 节点(要么单点要么没有过半保护)。
  • 加分项:「Controller 节点可以与 Broker 角色合并(教学环境),生产推荐分离」。
  • 易错点:以为 Quorum 越多越好——Quorum 写入要过半 ack,节点多反而慢。

6.6 KRaft 的元数据日志快照是怎么做的? ⭐⭐⭐⭐

  • 考察点:Raft 快照。
  • 标准答案
    1. __cluster_metadata 是无限增长的日志,不能无限保留。
    2. Active Controller 周期性把当前元数据状态写成快照文件(<offset>-<epoch>.checkpoint)。
    3. Broker 启动时先加载最新快照 + 应用快照后的增量日志,恢复元数据。
    4. 老于快照 offset 的日志可以截断。
  • 加分项:用 kafka-metadata-shell.sh --snapshot <file> 可以离线交互式查看。
  • 易错点:以为快照是 ZK 那种全量 dump。

6.7 ZK 时代 Controller 切换为什么慢?KRaft 怎么改进的? ⭐⭐⭐

  • 考察点:扩展性瓶颈。
  • 标准答案
    • ZK 时代:Controller 切换 → 新 Controller 从 ZK 全量拉所有元数据 → 通过 RPC 广播给所有 Broker → 百万分区时几分钟。
    • KRaft 时代:新 Active Controller 已经在 Quorum 里同步过日志,无需全量拉;Broker 也在持续订阅日志,元数据本就是最新的,切换是 ms 级
  • 加分项:补「KRaft 之后大集群 / 多分区的扩展上限主要在 Broker 本身的句柄数 / 内存了,元数据不再是瓶颈」。
  • 易错点:以为 KRaft 仅仅是「换了存储后端」。

6.8 KRaft 集群里 node.idbroker.id 是什么关系? ⭐⭐

  • 考察点:配置变化。
  • 标准答案
    • broker.id(ZK 时代):Broker 唯一 ID。
    • node.id(KRaft 时代):进程唯一 ID(既能是 broker 也能是 controller,所以不叫 broker.id)。
    • KRaft 模式下 node.id 必填且唯一;不再用 broker.id
  • 加分项:「controller.quorum.voters 里的数字也是 node.id」。
  • 易错点:仍用 broker.id

七、EOS / 事务

关联章节:Ch 13 EOS 真相 / Ch 12 消息语义

7.1 At Most Once / At Least Once / Exactly Once 工程含义分别是什么? ⭐⭐⭐

  • 考察点:消息语义基础。
  • 标准答案
    • At Most Once(最多一次):可能丢,不会重复。例:acks=0 Producer + 自动提交 Consumer。
    • At Least Once(至少一次):不丢,可能重复。例:acks=all Producer + 处理后手动提交 Consumer,Producer 重试导致重复 / Consumer 处理后宕机未提交。
    • Exactly Once(精确一次):不丢不重复。需 Kafka EOS 协议 + 业务幂等。
  • 加分项:「实战中最常见的是 At Least Once + 业务幂等——比 EOS 简单,效果一致」。
  • 易错点:以为分布式系统中 Exactly Once「天然不可能」(Kafka 在 consume-process-produce 链路里确实做到了)。

7.2 幂等 Producer 与事务 Producer 的区别? ⭐⭐⭐⭐

  • 考察点:两层机制。
  • 标准答案
    • 幂等 Producer:单分区维度去重,靠 (PID, partition, seq)。不能跨分区原子写
    • 事务 Producer:跨分区原子写 + 重启后保 Producer 身份(PID Fencing)+ Consumer 端 read_committed 过滤未提交。
  • 加分项:「幂等是事务的基础(事务自动开幂等);只有跨分区或者重启场景才需要事务」。
  • 易错点:以为「开了幂等就 = EOS」。

7.3 事务的两阶段提交流程是怎样的? ⭐⭐⭐⭐

  • 考察点:协议细节。
  • 标准答案
    1. initTransactions():Producer 注册到 Transaction Coordinator(__transaction_state 某分区的 Leader 即对应 Coordinator),拿 PID + Epoch;老 epoch 的 Producer 立即被 fence。
    2. beginTransaction() → 业务 send 消息(带 PID + Epoch)+ sendOffsetsToTransaction(offsets, group)(Streams 的「消费 offset 也加入事务」)。
    3. commitTransaction():① Coordinator 写 PREPARE_COMMIT 到 __transaction_state;② 向所有参与分区的 Leader 写控制消息(commit marker);③ 写 COMPLETE_COMMIT。任一步失败回滚。
    4. Consumer read_committed:跳过未提交事务的消息(看到 commit / abort marker 后才决定是否消费)。
  • 加分项:abort 路径类似,控制消息是 abort marker。
  • 易错点:以为事务是单 Broker 内的事;其实是跨分区跨 Broker 的两阶段。

7.4 PID Fencing 是什么?为什么需要? ⭐⭐⭐⭐

  • 考察点:进程级唯一性保证。
  • 标准答案
    1. Producer 重启后用同 transactional.id 重新 initTransactions() → Coordinator 给它一个新 epoch(递增)。
    2. 老进程(如果还活着或者网络分区里挣扎)拿旧 epoch 发消息会被 Broker / Coordinator 拒绝(ProducerFencedException)。
    3. 这就保证了任意时刻同一 transactional.id 只有一个 Producer 实例在写
  • 加分项:「这是『同一服务多实例不能共用 transactional.id』的根因」。
  • 易错点:以为 Fencing 是按 PID 而非 (transactional.id, epoch)。

7.5 Consumer 端 isolation.level=read_committed 怎么过滤未提交消息? ⭐⭐⭐⭐

  • 考察点:消费侧 EOS。
  • 标准答案
    1. Broker 在分区内同时存储「正常消息」+「事务控制消息(commit / abort marker)」。
    2. read_committed 时,Consumer 看到一个事务的 marker 之前buffer该事务的消息;marker = commit → 释放消息给业务;marker = abort → 丢弃。
    3. LSO(Last Stable Offset):消费者最多能读到的 offset = 当前所有「未决事务」中最早一个开始的 offset 之前。
  • 加分项:「这就是 read_committed 比 read_uncommitted 略慢的原因——要等 marker」。
  • 易错点:以为 Broker 物理删除了 abort 消息(其实保留,靠 LSO 过滤)。

7.6 consume-process-produce 链路怎么实现 EOS? ⭐⭐⭐⭐

  • 考察点:Streams EOS 核心模式。
  • 标准答案
    1. Producer 开事务。
    2. Consumer poll 拿一批 record。
    3. 业务处理(可能 send 新消息到下游 Topic)。
    4. producer.sendOffsetsToTransaction(consumed-offsets, group-id):把消费的 offset 也作为事务的一部分提交。
    5. producer.commitTransaction():原子提交「下游消息 + 消费 offset」。
  • 加分项:失败重试时整批被 abort,下次从原 offset 重新 poll → 不重复 / 不丢。
  • 易错点:忘了 sendOffsetsToTransaction,导致 offset 提交在事务外(可能重消费)。

7.7 EOS 性能代价有多大?什么时候不该用? ⭐⭐⭐

  • 考察点:成本意识。
  • 标准答案
    1. 写延迟:每个事务有 PREPARE / COMMIT 两阶段,加上控制消息额外 IO,单事务延迟 +几十 ms ~ 几百 ms。
    2. 吞吐:事务越频繁吞吐越低;一条一事务的极端用法吞吐能掉 10 倍。
    3. 不该用:① 业务能幂等去重;② 单分区写且不需要跨实例 Fencing;③ 高 QPS 的指标 / 日志类。
  • 加分项:「Streams EOS 通过『一个 task 的所有 send 走一个事务,commit 间隔 100ms ~ 几秒』摊薄成本」。
  • 易错点:滥用事务做小消息。

7.8 read_committed 消费时 Lag 看起来比实际多怎么回事? ⭐⭐⭐

  • 考察点:LSO 与 Lag 的关系。
  • 标准答案
    • kafka-consumer-groups --describe 中的 Lag = LogEndOffset - CURRENT-OFFSET,但 read_committed 实际能消费到的最大 offset 是 LSO(≤ HW ≤ LEO)。
    • 当某个事务长时间未 commit/abort(hang 事务),LSO 不会推进,看起来 Lag 一直涨,但其实是「等事务 marker」。
  • 加分项:监控 LastStableOffset 与 LEO 的差,如果持续大说明有未决事务。
  • 易错点:误以为消费跑得慢,其实是事务卡住。

八、性能与底层

关联章节:Ch 08 高吞吐底层 / Ch 04 Producer / Ch 19 运维

8.1 Kafka 高吞吐的「五大基石」分别是什么? ⭐⭐⭐⭐

  • 考察点:性能根源。
  • 标准答案
    1. 顺序写:Partition 是追加日志,避免随机寻道。
    2. PageCache:Broker 不用 JVM 堆做缓存,让 OS 用空闲内存做页缓存。
    3. 零拷贝(sendfile):Consumer fetch 直接从 PageCache → socket,避免用户态。
    4. 批量 + 压缩:Producer 攒批 + 压缩,端到端节省 CPU / 网络 / 磁盘。
    5. mmap 索引:稀疏索引文件 mmap 后查找 = 内存访问。
  • 加分项:还可以补「分区并行 + 客户端拉模型 + 副本只走 Leader」三点辅助。
  • 易错点:只说零拷贝。

8.2 PageCache 在哪些场景失效? ⭐⭐⭐

  • 考察点:性能洞察。
  • 标准答案
    1. Cold Read:读很久之前的数据 → 不在 PageCache → 走磁盘。
    2. 大量随机读:消费者 seek 到不同位点反复跳。
    3. PageCache 被竞争:同机部署其他进程占内存(典型:Broker 共部署 Flink)。
    4. 跨机房消费:网络 RTT 让 PageCache 命中也意义不大。
  • 加分项:监控 OS cat /proc/meminfoCachedActive(file),以及 Broker 的 BytesOutPerSec / MessagesInPerSec 比值。
  • 易错点:以为 PageCache 永远命中。

8.3 一个 Broker 单机能跑多少分区?瓶颈在哪? ⭐⭐⭐

  • 考察点:单机容量。
  • 标准答案
    • ZK 时代:单 Broker 4000 ~ 8000 分区(含副本)开始有压力,瓶颈在 Controller 广播、文件句柄、JVM 元数据。
    • KRaft 时代:可达数万 ~ 数十万,瓶颈转移到磁盘 IO、网络、内存。
  • 加分项:「单分区也有上限:每分区一组 file handle、segment buffer、副本同步线程;分区越多 PageCache 命中率越差(碎片化)」。
  • 易错点:报一个绝对数字(要看版本 + 硬件)。

8.4 Producer 端 linger.ms 设大有什么风险? ⭐⭐

  • 考察点:调参权衡。
  • 标准答案
    1. 端到端延迟变高(每条消息至少等 linger.ms)。
    2. buffer.memory 占用变多(攒得更久)→ 极端可能 buffer exhausted。
    3. Producer 关闭 / 应用 crash 时,buffer 里未发送的消息全丢。
  • 加分项:「优雅关闭时 producer.close(timeout) 会 flush buffer,但极端 SIGKILL / OOM 仍可能丢」。
  • 易错点:以为 linger.ms 越大越好。

8.5 怎么排查 Broker 端 IO 瓶颈? ⭐⭐⭐

  • 考察点:性能排障方法。
  • 标准答案
    1. OS 层 iostat -xm 1:看 %util / await / svctm / r_await / w_await
    2. Broker JMX:RequestHandlerAvgIdlePercent < 0.3 表示 IO 线程吃紧。
    3. 是 Produce 还是 Fetch 慢:TotalTimeMs,request=Produce vs Fetch 分桶。
    4. LogFlushRateAndTimeMs 突增 → 看是不是有外部 fsync 触发。
    5. 看 PageCache 命中:vmstat -spages paged in / out
  • 加分项:用 perf / bpftrace 看 sendfile 调用。
  • 易错点:只看 Broker JMX 不看 OS 层指标。

8.6 Broker 给多大堆?为什么不能太大? ⭐⭐⭐

  • 考察点:JVM 调优经验。
  • 标准答案
    • 推荐 6 ~ 10G,最多 16G
    • 太大:Full GC 单次几十秒卡住,触发 Controller 切换 + Rebalance + ISR 收缩 + 客户端 timeout。
    • 太小:元数据多 / 复杂客户端连接时 OOM。
  • 加分项:「OS 留出 ≥ 50% 内存给 PageCache 才是 Kafka 推荐内存模型」;4.0+ 推荐 generational ZGC。
  • 易错点:「越大越快」的旧观念。

8.7 网络成为瓶颈时怎么调? ⭐⭐⭐

  • 考察点:网络优化。
  • 标准答案
    1. num.network.threads / num.io.threads
    2. socket.send.buffer.bytes / socket.receive.buffer.bytes(高带宽链路)。
    3. Producer / Consumer 端开zstd 压缩减少字节量。
    4. Follower fetch 用 KIP-392 最近副本(跨 AZ 场景)。
    5. 副本 fetch 限流(reassignment 时 throttle 防打满)。
  • 加分项:补「跨机房用 MirrorMaker 2 异步复制,避免主链路 ack」。
  • 易错点:只调 send/recv buffer。

8.8 Kafka 的 transferTo 在什么情况下会退化? ⭐⭐⭐⭐

  • 考察点:底层细节。
  • 标准答案
    1. TLS 加密(SSL listener):消息要在用户态加密 → sendfile 退化为传统 read+write。
    2. 大消息超过单次 transferTo 上限:分多次拷贝。
    3. 数据不在 PageCache:要先从磁盘 read 再 sendfile。
  • 加分项:「这就是为什么 SSL 通信吞吐通常比 PLAINTEXT 低 30%~50%;想兼顾性能 + 加密,4.0 之后看 KIP-1041 等优化进展」。
  • 易错点:以为 sendfile 永远生效。

九、Streams / Connect / Schema Registry

关联章节:Ch 16 Connect / Ch 17 Streams / Ch 18 Schema Registry

9.1 Stream 与 Table 的「二元性」是什么? ⭐⭐⭐⭐

  • 考察点:流处理理论。
  • 标准答案
    1. Stream:无界事件序列(每条事件是 fact,不可变)。
    2. Table:某 Key 的最新状态(可变)。
    3. 二元性:把 Stream 按 Key 聚合 → Table(如「订单事件流 → 当前订单状态表」);把 Table 的变更看作 Changelog → Stream(如「用户表的变更 → 用户事件流,CDC 模式」)。
  • 加分项:补「Kafka Streams 用 KStream / KTable / GlobalKTable 三种抽象表达;KTable 自动用 changelog topic 持久化状态」。
  • 易错点:以为 Stream 和 Table 是两种 API,其实是同一份数据的两个视角。
  • 考察点:流处理选型。
  • 标准答案
    • Streams:Java 库,跑在业务进程里,无需独立集群;状态存本地 RocksDB + changelog Topic;适合轻量、与微服务一起部署的场景。
    • Flink:独立集群(JobManager + TaskManager),强大的窗口 / EventTime 支持、checkpoint 语义、SQL;适合重量级、跨数据源、ETL / 复杂分析
  • 加分项:「Streams + ksqlDB 在 Kafka 内部生态闭环;Flink 是更通用的流引擎」。
  • 易错点:以为 Streams 性能差于 Flink(同等场景吞吐相当)。

9.3 Kafka Connect 是什么?Standalone vs Distributed? ⭐⭐

  • 考察点:Connect 基础。
  • 标准答案
    • Connect = 把外部系统(DB / FS / S3 / ES / Kafka)跟 Kafka 串起来的标准框架;分 Source / Sink。
    • Standalone:单进程、配置文件本地、状态本地,开发 / 小规模用。
    • Distributed:集群、配置/offset/status 都存 Kafka 内部 Topic(connect-configs/offsets/status)、REST API 提交 connector,生产推荐
  • 加分项:补「Source 默认 At Least Once;EOS Source(KIP-618)需要支持事务」。
  • 易错点:把 Connect 当成 Streams 的一部分(其实是独立项目)。

9.4 SMT 是什么?举几个常用 SMT。 ⭐⭐⭐

  • 考察点:Connect 数据加工。
  • 标准答案
    • SMT = Single Message Transforms,挂在 connector 上对单条消息做轻量变换。
    • 常用:InsertField(注入字段)、ReplaceField(改名 / 删字段)、MaskField(脱敏)、Cast(类型转换)、TimestampConverter(时间转换)、RegexRouter(按正则改 Topic 名做路由)、Flatten(拍平嵌套)、Debezium 自带 ExtractNewRecordState(CDC 解包)。
  • 加分项:「重业务逻辑不该塞到 SMT 里,那是 Streams 的活儿;SMT 只做轻量变换」。
  • 易错点:把 SMT 当成 Streams 替代。

9.5 Schema Registry 的 BACKWARD / FORWARD / FULL 兼容性策略各自允许什么改动? ⭐⭐⭐⭐

  • 考察点:Schema 演进规则。
  • 标准答案
    • BACKWARD(默认):新 schema 能读老数据。允许:删字段、给字段加默认值。生产者先升级。
    • FORWARD老 schema 能读新数据。允许:加字段、删带默认值的字段。消费者先升级。
    • FULL:BACKWARD + FORWARD 都满足。允许:加 / 删带默认值的字段。
    • NONE:不校验,自由改(危险)。
    • _TRANSITIVE 后缀:与全部历史版本比较;不加则只与上一版本。
  • 加分项:「BACKWARD_TRANSITIVE 是大厂常用默认;新增必填字段是常见踩坑(任何策略都不允许)」。
  • 易错点:把 BACKWARD / FORWARD 方向记反。

9.6 Avro / Protobuf / JSON Schema 三选一怎么选? ⭐⭐⭐

  • 考察点:序列化选型。
  • 标准答案
    • Avro:Schema-First、紧凑二进制、跨语言;Confluent 生态首选;schema 演进规则成熟。
    • Protobuf:跨语言、字段 tag 友好演进、Google 生态;适合公司内已有 gRPC / Protobuf。
    • JSON Schema:可读、容易调试、字段多时变大;适合对人友好的场景。
  • 加分项:「与 Schema Registry 都能集成;Avro 是 Kafka 社区事实默认」。
  • 易错点:JSON 直接发不带 Schema → 没法演进 → 业务踩坑。

9.7 KStream 与 KTable 有什么区别?什么是 GlobalKTable? ⭐⭐⭐⭐

  • 考察点:Streams 抽象。
  • 标准答案
    • KStream:无状态事件流,每条 record 是独立事实。
    • KTable:有状态 K-V 表,同 Key 后写覆盖前写,背后由 compact changelog topic 持久化。
    • GlobalKTable:每个 Streams 实例全量加载一份到本地 RocksDB(不按分区切);适合「小维表 join」(避免大流分区与维表分区不一致带来的额外 repartition)。
  • 加分项:「KTable join 要求两侧 co-partition;GlobalKTable join 不要求」。
  • 易错点:把 GlobalKTable 用在大表上撑爆本地。

9.8 Kafka Streams 的 Time Semantics 有哪些? ⭐⭐⭐

  • 考察点:事件时间。
  • 标准答案
    • Event Time:消息内嵌业务时间戳。
    • Processing Time:消息被处理的墙上时间。
    • Ingestion Time:消息被 Broker 接收的时间。
    • default.timestamp.extractor 决定走哪种。
  • 加分项:窗口(Tumbling / Hopping / Sliding / Session)默认走 Event Time;处理乱序需 grace.ms
  • 易错点:把 Processing Time 用于业务统计 → 重启 / 消费慢时结果不一致。

9.9 Connect 的 Dead Letter Queue 怎么用? ⭐⭐⭐

  • 考察点:错误处理。
  • 标准答案
    • errors.tolerance=all + errors.deadletterqueue.topic.name=<dlq-topic> + errors.deadletterqueue.context.headers.enable=true
    • 反序列化失败 / Sink 写下游失败 → 原消息进 DLQ Topic,附带错误 Header(异常类、消息原 Topic、offset 等),不会卡住任务。
  • 加分项:DLQ Topic 副本数 ≥ 3,业务侧定时扫描修补。
  • 易错点:用了 DLQ 但没有人监控 DLQ → 数据无声丢失。

十、运维与排障

关联章节:Ch 19 可观测性与运维 / Ch 20 踩坑 / appendix_pitfalls.md

10.1 你看哪些 JMX 指标判断集群健康? ⭐⭐⭐⭐

  • 考察点:监控经验。
  • 标准答案
    1. 可用性OfflinePartitionsCount(>0 报警)、UnderReplicatedPartitions(持续 >0 报警)、UnderMinIsrPartitionCount(>0 影响写)。
    2. ControllerActiveControllerCount(集群和 = 1,否则异常)、LeaderElectionRateAndTimeMsUncleanLeaderElectionsPerSec(>0 报警)。
    3. 吞吐MessagesInPerSec / BytesInPerSec / BytesOutPerSec 趋势。
    4. 延迟TotalTimeMs,request=Produce/Fetch 各分位值。
    5. 资源RequestHandlerAvgIdlePercent / NetworkProcessorAvgIdlePercent(< 0.3 吃紧)。
    6. 客户端:Producer record-error-rate / buffer-exhausted-rate、Consumer records-lag-max / last-rebalance-seconds-ago
  • 加分项:用 Prometheus JMX Exporter + Grafana 现成面板;KMinion / kafka-exporter 工具。
  • 易错点:只看 server.log 不看 JMX。

10.2 一个分区 Lag 持续不下降,怎么排查? ⭐⭐⭐⭐

  • 考察点:综合排障。
  • 标准答案
    1. 看消费侧:是否在 Rebalance(last-rebalance-seconds-ago 频繁低)?是否单消费实例?业务处理时间多久?
    2. 看 Broker:该分区 Leader 是否 OK?UnderReplicated?磁盘 IO?
    3. 看 Topic:流量是否突增?分区是否热点?
    4. 客户端日志:是否有反序列化错?是否在 commit 失败?
    5. kafka-consumer-groups --describe --group g1 看 LAG 是否在涨,对比 OWNER 是否换。
  • 加分项:临时改 max.poll.records 减小批次、看是否能消化;或临时 scale-out 消费实例。
  • 易错点:只看消费侧 / 只看 Broker,不全面。

10.3 集群分区重分配(reassignment)怎么做?要注意什么? ⭐⭐⭐

  • 考察点:扩缩容操作。
  • 标准答案
    1. 准备 topics.json 列出要迁的 Topic + 要参与的 broker.list。
    2. kafka-reassign-partitions --generate 生成 current / proposed 两段 JSON。
    3. 人工审 proposed → --execute 启动迁移。
    4. 必须加 --throttle <bytes/s> 限速(如 50MB/s),否则副本同步打满网络影响线上。
    5. --verify 看进度;完成后用 kafka-leader-election --election-type preferred 切回 preferred Leader。
  • 加分项:用 LinkedIn 的 Cruise Control 自动平衡,比手工写 JSON 安全得多。
  • 易错点:不限流直接迁 → 网络打爆 → 业务超时。

10.4 你常调的 5 个 Broker / Producer / Consumer 关键参数是什么? ⭐⭐⭐

  • 考察点:经验深度。
  • 标准答案(参考答案,按场景给):
    • Brokermin.insync.replicaslog.retention.hourslog.segment.bytesnum.io.threadsunclean.leader.election.enable
    • Produceracksenable.idempotencelinger.ms + batch.sizecompression.typetransactional.id
    • Consumerenable.auto.commitmax.poll.recordsmax.poll.interval.mspartition.assignment.strategygroup.instance.id
  • 加分项:每条都能对应一个生产案例(参考 appendix_pitfalls.md)。
  • 易错点:把默认值当推荐值。

10.5 磁盘满了怎么处理? ⭐⭐⭐

  • 考察点:紧急处理。
  • 标准答案
    1. 绝对不要 rm log 文件
    2. kafka-configs --alter --add-config retention.ms=3600000 临时把某 Topic 保留期调小 → 几分钟自动清。
    3. 或用 kafka-delete-records.sh 删指定 offset 之前的数据。
    4. 急扩容:挂新盘到 log.dirs(KRaft 后支持热加)。
    5. 长期:评估容量;磁盘告警阈值 80%。
  • 加分项:「3 副本时只挂掉 1 个 Broker 仍能写,所以紧急情况下也可以杀掉满盘 Broker,让流量切到剩下 2 个,然后慢慢清理」。
  • 易错点:手动 rm → 索引文件不一致,数据无法恢复。

10.6 KRaft Controller 全挂了怎么办? ⭐⭐⭐⭐

  • 考察点:极端故障预案。
  • 标准答案
    1. 集群仍能继续读写已有 Topic(Broker 用本地缓存的元数据),但不能创建 / 删除 Topic、Leader 切换
    2. 立即恢复 Quorum 中至少过半节点(3 节点 Quorum 至少 2 个)→ 选出新 Active Controller。
    3. 如果半数以上 Controller 数据卷损坏 → 用最新快照 + 备份恢复。
  • 加分项:「这就是 Controller 数据卷必须独立、必须有快照备份的原因」。
  • 易错点:以为 Controller 挂了集群立刻死。

10.7 怎么定位某条消息「丢了」? ⭐⭐⭐⭐

  • 考察点:丢消息排查。
  • 标准答案
    1. 业务侧确认:业务 ID → Producer 端日志确认是否真的 send 成功(onSuccess 回调有 partition + offset)。
    2. 没 send 成功 → Producer 配置审:acks / retries / delivery.timeout.ms
    3. send 成功有 offset → kafka-console-consumer --offset <N> --partition <P> --max-messages 1 直接读看是否还在 Broker。
    4. Broker 上没了 → 是不是被 retention 清理 / Compaction 压掉?查 LogStartOffset
    5. 还在 → 消费侧问题:是否 commit 错位?是否有反序列化跳过?
  • 加分项:「acks=1 + Leader 切换 = 经典丢点」「自动提交 = 经典消费侧丢点」。
  • 易错点:只查消费侧不查 Producer 端。

10.8 跨机房双活 Kafka 怎么做? ⭐⭐⭐

  • 考察点:架构设计。
  • 标准答案
    • 方案 A:单集群拉伸跨机房(5 节点 Controller 跨 3 机房)→ 强一致但延迟高、网络分区风险。
    • 方案 B:双集群 + MirrorMaker 2 双向同步 → 各机房独立 Kafka,业务就近读写,MM2 异步复制;存在最终一致;冲突要业务侧处理。
    • 方案 C:Confluent Cluster Linking(商业版 / Pulsar Geo-Replication 风格)→ 平台级支持。
  • 加分项:「金融级强一致选 A 但接受写延迟;多数互联网选 B」。
  • 易错点:以为跨机房 = 拉伸单集群一定行。

十一、横向对比题

关联文档:appendix_kafka_vs_others.md

11.1 Kafka 与 RabbitMQ 的核心区别?分别适合什么场景? ⭐⭐⭐⭐

  • 考察点:经典选型。
  • 标准答案
    • 存储:Kafka 顺序日志;RabbitMQ 内存队列 + 持久化兜底。
    • 消费:Kafka 拉 + offset;RabbitMQ 推 + ACK 即删。
    • 路由:Kafka 仅 Key Hash;RabbitMQ Exchange + Routing Key 灵活。
    • 可重放:Kafka ✅;RabbitMQ ❌(除 Stream)。
    • 吞吐量级:Kafka GB/s;RabbitMQ 数十万 msg/s。
    • 延迟:Kafka 10ms+;RabbitMQ ms 级。
    • 选 Kafka:数据总线、流处理、长保留;选 RabbitMQ:业务异步、复杂路由、低延迟 RPC。
  • 加分项:很多公司两个都用:业务用 RabbitMQ,数据总线用 Kafka。
  • 易错点:把 RabbitMQ 当成「带持久化的 Kafka」(路由模型完全不同)。

11.2 Kafka 与 RocketMQ 的核心区别? ⭐⭐⭐⭐

  • 考察点:同源对比。
  • 标准答案
    • 同源:底层都是顺序追加日志(CommitLog ≈ segment)、都拉模型、都有 Group 消费。
    • 业务能力差异:RocketMQ 服务端 Tag/SQL 过滤、原生延迟消息(18 级别 + 5.x 任意秒)、半消息事务;Kafka 没这些(要客户端自实现)。
    • 生态差异:Kafka Streams / Connect / Schema Registry / 多语言 完整;RocketMQ Java 生态强、其他语言弱。
    • 元数据:Kafka KRaft;RocketMQ NameServer(无状态)。
    • 选 Kafka:数据总线 / 大数据 / 多语言;选 RocketMQ:业务消息 / 顺序 / 事务 / 延迟一站式 + Java 团队。
  • 加分项:RocketMQ 5.x 有 RocketMQ Streams、Proxy 模式等,向 Kafka 的生态靠拢。
  • 易错点:以为 RocketMQ 就是「中文版 Kafka」(业务能力差异巨大)。

11.3 Kafka 与 Pulsar 的核心区别? ⭐⭐⭐⭐

  • 考察点:架构对比。
  • 标准答案
    • 架构:Kafka 计算 + 存储耦合(Broker 自己存);Pulsar 分离(Broker 无状态 + BookKeeper 存储)。
    • 多租户:Pulsar 三级(Tenant / Namespace / Topic)原生强;Kafka 弱(靠 ACL + 命名约定)。
    • 跨地域:Pulsar Geo-Replication 内置;Kafka MirrorMaker 2。
    • 弹性:Pulsar 加 Broker 不需要数据迁移;Kafka 加 Broker 需要 reassignment。
    • 运维:Pulsar 三组件(Broker + Bookie + ZK/Oxia)比 Kafka 复杂。
    • 选 Pulsar:SaaS 多租户 / 跨地域 / 弹性强;选 Kafka:生态成熟 / 团队熟悉 / 运维简单。
  • 加分项:「Pulsar 端到端延迟略高于 Kafka(多一跳 Bookie)」。
  • 易错点:把「分离架构」当万能优点。

11.4 Kafka 与 Redis Stream 的边界? ⭐⭐⭐

  • 考察点:量级感知。
  • 标准答案
    • Redis Stream 适合:消息量小(< 百万 / 天)、亚毫秒延迟、已经在用 Redis、不需要长保留。
    • Kafka 适合:长期保留、多消费方共享、流处理、GB/s 吞吐、跨数据中心。
  • 加分项:「Redis Stream PEL ≈ Kafka 未提交 offset,但断电后受 AOF 持久化策略影响」。
  • 易错点:在 Redis Stream 上跑大数据流。

11.5 Kafka 与 NATS JetStream 的边界? ⭐⭐⭐

  • 考察点:云原生选型。
  • 标准答案
    • NATS JetStream 适合:K8s / 边缘 / 极简运维 / 多协议(NATS / MQTT / WebSocket)/ 极低延迟。
    • Kafka 适合:大数据生态丰富、长期保留、流处理上下游成熟。
  • 加分项:「NATS Subject 通配符服务端路由比 Kafka 灵活」。
  • 易错点:以为 NATS 适合大数据场景(生态薄)。

11.6 同样要做「秒杀削峰」,选 Kafka 还是 RocketMQ 还是 RabbitMQ? ⭐⭐⭐⭐

  • 考察点:场景化选型。
  • 标准答案
    • Kafka:吞吐高、能扛百万 QPS 入流,下游慢消费可堆积;缺点是延迟略高、不支持优先级和延迟。
    • RocketMQ:吞吐与 Kafka 同量级,自带顺序 / 事务 / 延迟消息最适合秒杀这种业务消息
    • RabbitMQ:吞吐到几十万 QPS 就开始吃力,但路由灵活,小规模业务首选
  • 加分项:「秒杀场景往往要『下单后 30 秒未付款关单』,这种延时关单 RocketMQ 一行配置;Kafka 要自实现轮转 Topic」。
  • 易错点:默认上 Kafka 不考虑业务诉求。

11.7 同样做「实时数仓」,Kafka 与 Pulsar 怎么选? ⭐⭐⭐

  • 考察点:实时数仓选型。
  • 标准答案
    • Kafka:与 Flink / Spark / Streams / ksqlDB / Connect 生态深度集成,案例最多;强烈推荐。
    • Pulsar:也支持 Flink Source / Sink,但生态稳定性 / 文档丰度略低;多租户场景下 Pulsar 占优。
  • 加分项:「数仓上下游基本都是『一份数据多吃』,Kafka 的 Consumer Group 模型天然契合」。
  • 易错点:用 RabbitMQ 做实时数仓(路由模型不匹配)。

11.8 想引入消息中间件做「业务异步解耦」,团队对 MQ 完全没经验,选什么? ⭐⭐⭐

  • 考察点:低门槛选型。
  • 标准答案
    • 优先 RabbitMQ(管理台直观、路由灵活、文档多),中文资料齐全。
    • 团队 Java 较强且业务有顺序 / 事务 / 延迟需求 → RocketMQ
    • 已经用了 Redis 且消息量小 → Redis Stream
    • 数据驱动 / 流处理需求 → Kafka
  • 加分项:「上 Kafka 需要更多『分区设计 / Rebalance / EOS』踩坑成本,团队成本评估」。
  • 易错点:盲目跟风上 Kafka。

11.9 「日志型 MQ」一定全局有序吗?怎么权衡? ⭐⭐⭐

  • 考察点:顺序保证 vs 并行。
  • 标准答案
    • 公理:日志型 MQ「全局有序」 = 单分区 / 单队列 / 单 subject,牺牲并行度
    • 实战推荐:业务 Key 维度有序(如同一订单同一用户),用 Key Hash 落同分区;不追求全局。
  • 加分项:「需要全局有序的极少数场景(如金融对账),可考虑用单分区 Topic + 单消费实例 + 加大批次」。
  • 易错点:以为「能开关」配一下就全局有序。

11.10 一个公司同时引入 Kafka + RabbitMQ 是不是冗余? ⭐⭐⭐

  • 考察点:架构合理性。
  • 标准答案
    • 不冗余。两者解决不同问题:RabbitMQ 是业务异步 / 复杂路由 / RPC 风格;Kafka 是数据总线 / 流处理 / 历史回放。
    • 业界常见拓扑:业务系统用 RabbitMQ 互相调用;同时把业务事件落到 Kafka 做数仓 + 风控 + 实时大屏。
  • 加分项:「事实上很多大厂同时维护 Kafka / RocketMQ / RabbitMQ 三套 MQ,每个解决一类问题」。
  • 易错点:「中间件越少越好」的极端思维。

🔥 高频题 Top 30(大厂一面必刷)

按出现频率从高到低(综合 200+ 份面经统计的主观经验排序)。建议连续 3 天每天刷一遍直到能 5 分钟内说完。

Rank题目难度出处
1Kafka 为什么快?(顺序写 + PageCache + 零拷贝 + 批量 + 压缩五大基石)⭐⭐⭐⭐第八章 / § 8.1
2Topic / Partition / Offset / Group 的关系§ 1.1
3acks=0/1/all 的语义差别⭐⭐⭐§ 3.1
4幂等 Producer + 事务 Producer 实现 EOS 的全链路⭐⭐⭐⭐§ 7.3 / § 3.6
5ISR / HW / LEO / Leader Epoch 一句话各是什么⭐⭐⭐⭐§ 5.1 / § 5.2 / § 5.3
6min.insync.replicas + acks=all 保证了什么⭐⭐⭐⭐§ 5.5
7unclean.leader.election 开了会怎样⭐⭐⭐⭐§ 5.4
8Rebalance 触发条件 + 4 种分配策略⭐⭐⭐⭐§ 4.4 / § 4.5
9Cooperative Rebalance vs Eager Rebalance⭐⭐⭐⭐§ 4.6
10KRaft 与 ZK 时代有什么区别⭐⭐⭐⭐§ 6.2
11自动提交 vs 手动提交,怎么避免丢消息⭐⭐⭐⭐§ 4.3
12max.poll.interval.ms 引发的 Rebalance 风暴⭐⭐⭐⭐§ 10.2 / 附录 C 4.1
13分区数选少 / 选多有什么后果⭐⭐⭐§ 1.4 / 附录 C 1.1
14Kafka vs RabbitMQ vs RocketMQ vs Pulsar⭐⭐⭐⭐§ 11 全章
15零拷贝 sendfile 省掉了什么⭐⭐⭐⭐§ 2.6
16PageCache 在 Kafka 里的作用⭐⭐⭐⭐§ 2.5
17Log Compaction 工作原理 + 必用场景⭐⭐⭐⭐§ 2.8
18一条消息从 Producer 到 Consumer 的完整链路⭐⭐⭐§ 1.3
19副本同步是 Leader 推还是 Follower 拉?为什么?⭐⭐⭐§ 5.6
20Static Membership 解决了什么问题⭐⭐⭐§ 4.7
21Schema Registry 兼容性策略 BACKWARD/FORWARD/FULL⭐⭐⭐⭐§ 9.5
22Stream 与 Table 二元性⭐⭐⭐⭐§ 9.1
23transactional.id 重复会怎样(PID Fencing)⭐⭐⭐⭐§ 7.4 / 附录 C 7.1
24一个分区 Lag 不降怎么排查⭐⭐⭐⭐§ 10.2
255 种压缩算法选哪个⭐⭐⭐§ 3.5
26默认分区器是什么?Sticky 有什么坑⭐⭐⭐§ 3.4 / 附录 C 2.4
27怎么定位「消息丢了」⭐⭐⭐⭐§ 10.7
28EOS 性能代价 / 什么时候不该用⭐⭐⭐§ 7.7
29跨机房双活 Kafka 怎么做⭐⭐⭐§ 10.8
30Kafka 单 Broker 能跑多少分区?瓶颈在哪⭐⭐⭐§ 8.3

🏃 面试 1 周冲刺路线

适合「下周一面,今天周日」的紧急救援。每天 3 ~ 4 小时。

Day 1(周日)· 把 Top 30 过一遍

  • 上午:Top 30 第 1 ~ 15 题,每题朗读 + 关键词手抄。
  • 下午:Top 30 第 16 ~ 30 题。
  • 晚上:把第 1 / 4 / 5 / 8 / 10 题(架构、EOS、ISR、Rebalance、KRaft)反复说 5 遍,能在 3 分钟内讲清楚为止。

Day 2(周一)· 架构 + 基础概念 + 存储

Day 3(周二)· Producer + Consumer + Rebalance

Day 4(周三)· 副本 + Controller / KRaft

Day 5(周四)· EOS + 性能

Day 6(周五)· 生态 + 运维 + 横向对比

Day 7(周六)· 模拟面试 + 查缺补漏

  • 上午:找一份大厂真题(搜「Kafka 面试题汇总」),自己掐表答。
  • 下午:把回答不流畅的题目对回本文档定位章节,重新精读 + 默背。
  • 晚上:早睡。明天加油!

📚 配套学习材料


💪 祝你拿到心仪的 offer!如果某道题的答案在面试时被问到细节,记得回到对应章节正文里有可运行的代码与图解可以反复消化。