Kafka 面试笔记
1 Kafka 是什么
Apache Kafka 是一种发布订阅消息系统,是一个分布式、分区的和可重复的日志服务。它是一个高吞吐量、分布式、支持分区(partition)、多副本(replica)、基于 Zookeeper 协调的实时消息系统。
2 与传统消息系统的区别
传统消息传递方式:
- 排队:一组用户从服务器读消息,每条消息只发给其中一个人;
- 发布-订阅:消息被广播给所有用户。
Kafka 与传统消息系统的三个关键区别:
- Kafka 持久化日志,可被重复读取并无限期保留;
- Kafka 是分布式系统,以集群方式运行、可灵活伸缩,内部通过复制数据提升容错和高可用;
- 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 检测、分布式同步、配置管理、识别新节点何时离开或连接集群、节点实时状态等。
判断节点是否存活的两个条件
- 节点必须能维护与 ZooKeeper 的连接(心跳机制检查);
- 如果节点是 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 高吞吐的实现
- 页缓存(Page Cache):读写文件依赖 OS 的页缓存,而不是 JVM 内部缓存,内存利用率高;
- sendfile 零拷贝:避免传统网络 IO 的四步流程,通过内核空间保存副本,字节从 socket 转移到磁盘;
- 顺序 IO + 常量时间 get/put;
- 支持 End-to-End 压缩;
- Partition 可很好横向扩展、提供高并发处理。
11 高效文件存储设计
- 把 topic 中一个 partition 的大文件分成多个小文件段,方便定期清理已消费完的文件,减少磁盘占用;
- 通过索引信息快速定位 message 和确定 response 的最大大小;
- 将 index 元数据全部映射到 memory,避免 segment file 的磁盘 IO;
- 索引文件稀疏存储,大幅降低 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 消息队列的作用
- 解耦:在系统中插入基于数据的接口层,两边只实现接口约束,可独立扩展修改;
- 顺序保证:队列本身有序,Kafka 保证一个 Partition 内消息有序;
- 异步通信:消息入队后不立即处理,在需要时再处理;
- 冗余:数据持久化直到被完全处理,规避丢失风险;
- 扩展性:解耦后增大处理频率只需增加处理过程;
- 峰值处理能力:消息队列使关键组件顶住突发访问压力,不必为峰值流量随时待命;
- 可恢复性:部分组件失效不影响整个系统,恢复后可继续处理队列中的消息;
- 缓冲:通过缓冲层帮助任务高效执行,控制数据流经过系统的速度。
15 其他 FAQ
- broker 意义:在 Kafka 集群中指服务器;
- 消息最大大小:Kafka 服务器可接收的消息最大 100 万字节(1M);
- producer 发送到 leader:producer 直接将数据发送到 broker 的 leader,所有节点会告知哪些节点活跃、目标 topic 分区 leader 在哪;
- 复制目的:确保已发布的消息不丢失,可在机器错误、程序错误或软件升级时使用;
- 副本长期在 ISR 中:说明该跟踪器无法像 leader 收集数据那样快速获取数据;
- 首选副本不在 ISR 中:控制器无法把 leadership 转移到首选副本;
- 生产后消息偏移:作为用户可以从 broker 获取补偿,SimpleConsumer 会获取包含偏移量作为列表的 MultiFetchResponse 对象。
阅读 —
·
全站 —