Kafka 面试笔记

Kafka 面试笔记

参考:Kafka 高效文件存储设计、Kafka 原理

1 Kafka 是什么

Apache Kafka 是一种发布订阅消息系统,是一个分布式、分区的和可重复的日志服务。它是一个高吞吐量、分布式、支持分区(partition)、多副本(replica)、基于 Zookeeper 协调的实时消息系统。

2 与传统消息系统的区别

传统消息传递方式:

  • 排队:一组用户从服务器读消息,每条消息只发给其中一个人;
  • 发布-订阅:消息被广播给所有用户。

Kafka 与传统消息系统的三个关键区别:

  1. Kafka 持久化日志,可被重复读取并无限期保留;
  2. Kafka 是分布式系统,以集群方式运行、可灵活伸缩,内部通过复制数据提升容错和高可用;
  3. Kafka 支持实时的流式处理。

3 Kafka 的优势

  • 高吞吐、低延迟:每秒可处理几十万条消息,延迟最低几毫秒;每个 topic 可分多个 partition,consumer group 对 partition 进行消费;
  • 可扩展性:集群支持热扩展;
  • 持久性、可靠性:消息持久化到本地磁盘,支持数据备份防止丢失;
  • 容错性:允许集群中节点失败(副本数为 n 则允许 n-1 个节点失败);
  • 高并发:支持数千个客户端同时读写。

4 核心概念

术语 含义
Topic 划分 Message 的逻辑概念,一个 Topic 可分布在多个 Broker 上
Partition Kafka 横向扩展和并行化的基础,每个 Topic 至少 1 个 Partition
Offset 消息在 Partition 中的编号,编号顺序不跨 Partition
Producer 向 Broker 发送/生产 Message
Consumer 从 Broker 取出/消费 Message
Broker 接受 Producer/Consumer 请求并把消息持久化到本地磁盘的服务器
Replication 以 Partition 为单位对消息冗余备份,每个 Partition 可配置至少 1 个副本
Leader 每个 Partition 的副本集合中选出的唯一副本,所有读写请求都由它处理,其他 Replicas 同步它的数据(类似 MySQL Binlog 同步)
ISR Replicas 的子集,表示当前 Alive 且能与 Leader “Catch-up” 的副本集合;从 Leader 拉取数据的延迟或条数超过阈值会被踢出 ISR
Controller 每个集群选举一个 Broker 担任,负责 Partition 的 Leader 选举、协调 Partition 迁移

5 ZooKeeper 的作用

Zookeeper 是一个分布式、高性能的协调服务。Kafka 不能越过 Zookeeper 直接联系 broker,一旦 Zookeeper 停止工作,就无法服务客户端请求。

作用:

  • 用于集群中不同节点之间的通信;
  • 提交偏移量(offset),节点失败时可以从之前提交的偏移量恢复;
  • leader 检测、分布式同步、配置管理、识别新节点何时离开或连接集群、节点实时状态等。

判断节点是否存活的两个条件

  1. 节点必须能维护与 ZooKeeper 的连接(心跳机制检查);
  2. 如果节点是 follower,必须能及时同步 leader 的写操作,延时不能太久。

6 消息传递语义

  • 最多一次:消息不会被重复发送,最多传输一次,但可能一次都不传输;
  • 最少一次:消息不会被漏发送,最少传输一次,但可能被重复传输;
  • 精确一次(Exactly once):不丢不重,每个消息传输一次且仅一次。

7 Pull 模式 vs Push 模式

Kafka 采用 producer 推送、consumer 拉取的 Pull 模式:

  • Push 缺点:broker 决定推送速率,对消费速率不同的 consumer 不好处理;当 broker 推送速率远大于 consumer 消费速率时 consumer 可能崩溃(如 Scribe、Apache Flume);
  • Pull 优点:consumer 可以自主决定消费速度,可以自主决定是否批量拉取;
  • Pull 缺点:broker 无消息时 consumer 会一直轮询。Kafka 通过参数让 consumer 阻塞直到新消息到达或消息数量达到特定量,实现批量发送。

8 数据一致性(Exactly Once 语义)

  • 生产端避免重复:每个分区使用单独写入器,遇到网络错误时检查该分区最后一条消息确认上次写入是否成功;或在消息中包含主键(UUID),在消费端做去重;
  • 消费端避免重复:消费端根据业务做反复制(幂等)。

9 ack 机制(request.required.acks)

值 行为 风险
0 生产者不等待 broker 的 ack 延迟最低、存储保证最弱,server 挂掉就丢数据
1 等待 leader 副本确认收到后发送 ack leader 挂掉且复制未完成,新 leader 也会导致数据丢失
-1 等待所有 follower 副本收到数据后才收到 leader 的 ack 数据不会丢失

10 高吞吐的实现

  1. 页缓存(Page Cache):读写文件依赖 OS 的页缓存,而不是 JVM 内部缓存,内存利用率高;
  2. sendfile 零拷贝:避免传统网络 IO 的四步流程,通过内核空间保存副本,字节从 socket 转移到磁盘;
  3. 顺序 IO + 常量时间 get/put;
  4. 支持 End-to-End 压缩;
  5. Partition 可很好横向扩展、提供高并发处理。

11 高效文件存储设计

  1. 把 topic 中一个 partition 的大文件分成多个小文件段,方便定期清理已消费完的文件,减少磁盘占用;
  2. 通过索引信息快速定位 message 和确定 response 的最大大小;
  3. 将 index 元数据全部映射到 memory,避免 segment file 的磁盘 IO;
  4. 索引文件稀疏存储,大幅降低 index 文件元数据占用空间。

Partition 数据如何保存到硬盘

  • topic 中多个 partition 以文件夹形式保存到 broker,每个分区序号从 0 递增,且消息有序;
  • Partition 文件下有多个 segment(xxx.index、xxx.log),segment 文件大小默认为 1G,超过则滚动新建一个 segment,并以上一个 segment 最后一条消息的偏移量命名。

消息格式

消息由固定长度的头部和可变长度的字节数组组成:消息长度 4 bytes + 版本号 1 byte + CRC32 校验码 4 bytes + 具体消息 n bytes。

12 分区与副本的放置

  • 副本因子不能大于 Broker 个数;
  • 第一个分区(编号 0)的第一个副本位置随机从 brokerList 选择,其他分区的第一个副本依次往后移;
  • 剩余副本相对第一个副本的位置由随机产生的 nextReplicaShift 决定。

新建分区在哪个目录创建

  • log.dirs 配置多个目录时,Kafka 会在含有分区目录最少的文件夹中创建新的分区目录(分区目录名为 Topic名+分区ID);
  • 注意是分区文件夹总数最少,而非磁盘使用量最少。给 log.dirs 新增磁盘后,新分区会先在这个新磁盘上创建,直到它不再是分区目录最少的目录。

13 消费者

  • 消费者消费时向 broker 发出 fetch 请求消费特定分区的消息,指定 offset 后消费该位置开始的消息,可回滚重新消费;
  • 消费者每次消费会记录物理偏移量(offset),下次接着上次位置继续消费;
  • 负载均衡:一个消费者组中一个分区对应一个消费者成员;组内有序、组间无序;组中成员太多会有空闲成员。

14 消息队列的作用

  1. 解耦:在系统中插入基于数据的接口层,两边只实现接口约束,可独立扩展修改;
  2. 顺序保证:队列本身有序,Kafka 保证一个 Partition 内消息有序;
  3. 异步通信:消息入队后不立即处理,在需要时再处理;
  4. 冗余:数据持久化直到被完全处理,规避丢失风险;
  5. 扩展性:解耦后增大处理频率只需增加处理过程;
  6. 峰值处理能力:消息队列使关键组件顶住突发访问压力,不必为峰值流量随时待命;
  7. 可恢复性:部分组件失效不影响整个系统,恢复后可继续处理队列中的消息;
  8. 缓冲:通过缓冲层帮助任务高效执行,控制数据流经过系统的速度。

15 其他 FAQ

  • broker 意义:在 Kafka 集群中指服务器;
  • 消息最大大小:Kafka 服务器可接收的消息最大 100 万字节(1M);
  • producer 发送到 leader:producer 直接将数据发送到 broker 的 leader,所有节点会告知哪些节点活跃、目标 topic 分区 leader 在哪;
  • 复制目的:确保已发布的消息不丢失,可在机器错误、程序错误或软件升级时使用;
  • 副本长期在 ISR 中:说明该跟踪器无法像 leader 收集数据那样快速获取数据;
  • 首选副本不在 ISR 中:控制器无法把 leadership 转移到首选副本;
  • 生产后消息偏移:作为用户可以从 broker 获取补偿,SimpleConsumer 会获取包含偏移量作为列表的 MultiFetchResponse 对象。
阅读 — · 全站 —
🎸 我的歌单 0 首