一句话结论
Kafka 是分布式日志系统——消息按序追加到 Partition 末尾(顺序写),消费者按 Offset 拉取(Pull 模型)。高吞吐靠顺序写 + 零拷贝 + 批量压缩。
核心概念
消息可靠性
Producer → Broker:
acks=0: 不等确认(可能丢)
acks=1: Leader 确认即可(Leader 宕机可能丢)
acks=all: 所有 ISR 确认(最安全,默认 min.insync.replicas≥2)
Broker 存储:
Partition 多副本 → ISR 机制
Consumer:
手动提交 Offset → 处理完再提交
高吞吐原理
顺序写
Kafka 的消息是追加到 Partition 文件末尾(append-only),
是顺序写入而非随机写入——磁盘顺序写速度可达 600MB/s+,
媲美内存随机写速度。数据不删不改,天然适合日志语义。
Page Cache + 零拷贝
Producer → Broker: 消息写入 Page Cache → 异步刷盘
Consumer ← Broker: sendfile() 零拷贝 → 数据从 Page Cache 直接到网卡
不经过用户态内存拷贝!
批量 + 压缩
Producer 批量发送(batch.size + linger.ms)
Broker 批量存储(Segment 文件)
Consumer 批量拉取(fetch.min.bytes)
消息支持 gzip/snappy/lz4/zstd 压缩
Partition 分配策略
Consumer Group 内的 Partition 分配:
Partition 1 ──→ Consumer A
Partition 2 ──→ Consumer B
Partition 3 ──→ Consumer C
Partition 4 ──→ Consumer A (一个 Consumer 可消费多个 Partition)
规则: 同组内一个 Partition 只能被一个 Consumer 消费
一个 Consumer 可以消费多个 Partition
Consumer 数量 > Partition 数量 → 多余的 Consumer 空闲
ISR(In-Sync Replicas)
保持与 Leader 同步的副本集合。副本落后超过 replica.lag.time.max.ms(默认 30s)被踢出 ISR。Leader 宕机 → Controller 从 ISR 中选新 Leader。
ISR 机制:
Partition 0: Leader(Broker1) Follower(Broker2) Follower(Broker3)
[0,1,2,3,4,5] [0,1,2,3,4,5] [0,1,2,3] ← 落后了!
↑ 超过阈值踢出 ISR
acks=all + min.insync.replicas=2:
至少 Leader + 1 个 Follower 确认才算写入成功
Rebalance 详细过程
触发条件:
- Consumer 加入或离开 Consumer Group
- Topic Partition 增加
- Consumer 心跳超时(session.timeout.ms 默认 45s)
Rebalance 流程(Eager 协议,Kafka 2.4+ 支持 Cooperative 增量重分配):
1. Group Coordinator 收到触发事件
2. 所有 Consumer 停止消费("Stop the World")
3. Coordinator 根据 partition.assignment.strategy 重新分配
4. Consumer 收到新分配,恢复消费
Eager vs Cooperative:
Eager: 全停 → 重新分配 → 全恢复(有停顿)
Cooperative: 只调整有变化的 Partition,其他继续消费(无停顿)
Offset 管理
Offset 提交:
自动提交(enable.auto.commit=true,默认 5s 间隔):
✅ 简单
❌ 可能丢消息(poll 后自动提交但还没处理完)
❌ 可能重复消费(处理完但还没到下一个提交间隔就崩溃)
手动提交:
✅ 精确控制(处理完再提交)
❌ 代码复杂
⚠️ 注意: 提交的是下次要消费的 offset(不是已消费的最后一条)
__consumer_offsets Topic:
Kafka 内部用这个 Topic 存 Consumer Group 的 Offset
(Kafka 0.9 之前 offset 存 Zookeeper → 性能瓶颈)
消息语义
幂等 Producer(enable.idempotence=true):
每条消息带 ProducerID + SequenceNumber
Broker 去重(同一 PID 的同一 SN 只持久化一次)
事务 Producer:
beginTransaction → 发送多条消息 → commitTransaction
原子性保证(所有消息要么都可见要么都不可见)
与 RabbitMQ 对比
高频面试问题
Q: Kafka 为什么吞吐这么高?
30 秒回答: 四个核心设计:① 顺序写(磁盘顺序写接近内存速度)② Page Cache + sendfile 零拷贝(数据从磁盘到网卡不经过用户态)③ 批量压缩(Producer/Consumer/Broker 三端批量)④ Partition 并行(水平扩展无上限)。
Q: Kafka 的 ISR 和 Raft 的 Quorum 有什么区别?
30 秒回答: ISR 是动态集合(跟不上就踢出),Raft Quorum 是固定多数派(2N+1 个节点过半数确认)。ISR 追求低延迟(等最快的几个),Raft 追求安全性(等多数派里的多数)。
Q: Rebalance 期间能消费吗?
30 秒回答: Eager Rebalance(老方案):不能,全部暂停。Cooperative Rebalance(Kafka 2.4+):可以,只暂停被重新分配的那部分 Partition。
项目相关
目前项目用 RabbitMQ。如果未来做设备日志采集、Agent 事件溯源会考虑 Kafka。
速记
Kafka=分布式日志+顺序写+零拷贝。Partition=分区有序。Consumer Group=分摊消费。ISR=同步副本。acks=all最安全。Rebalance=重新分配暂停消费。vs RabbitMQ=业务任务用RMQ,日志流用Kafka。