Kafka 深度解析
#组件 · #Kafka · #消息队列 · #流处理 · #发布订阅
从存储模型到高吞吐设计,从 Producer/Consumer 到可靠性保证,全面理解 Kafka 的核心思想。
一、架构概览
mermaid
flowchart TB
subgraph Cluster["Kafka 集群"]
subgraph B1["Broker 1"]
P0["Topic-A Partition 0<br/>(Leader)"]
R1["Topic-A Replica 1"]
R2["Topic-A Replica 2"]
end
subgraph B2["Broker 2"]
P1["Topic-A Partition 1<br/>(Leader)"]
R0a["Topic-A Replica 0"]
R2a["Topic-A Replica 2"]
end
subgraph B3["Broker 3"]
P2["Topic-A Partition 2<br/>(Leader)"]
R0b["Topic-A Replica 0"]
R1b["Topic-A Replica 1"]
end
end
ZK["ZooKeeper / KRaft<br/>元数据管理<br/>(Controller选举、分区分配)"]
ZK -->|"协调"| Cluster
Producer["Producer<br/>生产者"] -->|"写入消息"| Cluster
Cluster -->|"拉取消费"| Consumer["Consumer Group<br/>消费者组<br/>每个分区仅一个消费者"]二、存储模型(Kafka 高性能的核心)
2.1 分区与日志段
Topic: user-events (3 分区, 2 副本)
Partition 0 的文件布局:
/var/lib/kafka/data/user-events-0/
├── 00000000000000000000.log ← 日志段 (存储消息)
├── 00000000000000000000.index ← 稀疏索引 (offset → 文件位置)
├── 00000000000000000000.timeindex← 时间索引
├── 00000000000000102456.log ← 下一个日志段
├── 00000000000000102456.index
└── leader-epoch-checkpoint2.2 顺序写入 + Page Cache
Producer → Broker:
1. 追加消息到 .log 文件末尾 (顺序写, 仅追加)
2. 不立即 fsync, 依赖 OS Page Cache
3. 消费者读取时大概率命中 Page Cache (零磁盘 I/O)为什么顺序写这么快?
| 对比 | 随机写 | 顺序写 |
|---|---|---|
| 磁盘寻道 | 每次都要磁头移动 | 几乎没有 |
| 机械硬盘 | 100 IOPS | 6000+ IOPS |
| SSD | 也比顺序慢很多 | 极快 |
2.3 零拷贝 (Zero-Copy)
传统方式 (4 次拷贝 4 次上下文切换):
磁盘 → Read Buffer → 用户态 Buffer → Socket Buffer → 网卡
Kafka sendfile (2 次拷贝 2 次切换):
磁盘 → Read Buffer ────────────────→ Socket Buffer → 网卡
(DMA 直接传输)
节省: 用户态拷贝 + 2 次上下文切换2.4 日志清理策略
| 策略 | 参数 | 行为 |
|---|---|---|
| 删除 (默认) | cleanup.policy=delete | 按时间(retention.ms)或大小(retention.bytes)删除旧段 |
| 压缩 | cleanup.policy=compact | 同 Key 只保留最新 Value (适合 CDC/快照) |
| 两者 | cleanup.policy=compact,delete | 同时生效 |
日志压缩原理:
压缩前:
Key: A → v1 压缩后:
Key: B → v1 Key: A → v3
Key: A → v2 → Key: B → v2
Key: B → v2 Key: C → v1
Key: A → v3三、生产者
3.1 发送流程
Producer 端:
┌──────────┐ ┌──────────┐ ┌──────────┐
│序列化 │→ │分区器 │→ │消息累加器 │→ Sender 线程异步发送
│Serializer│ │Partitioner│ │RecordAccum│ (批量, 按 Broker 分组)
└──────────┘ └──────────┘ └──────────┘
│
┌─────┴─────┐
│ BufferPool │ ← 内存池复用 ByteBuffer
└───────────┘3.2 分区策略
java
// 默认: 指定 key → hash(key) % partition_count
// 未指定 → sticky (2.4+) 随机选分区攒批
// 自定义分区器
public class MyPartitioner implements Partitioner {
public int partition(String topic, Object key, byte[] keyBytes,
Object value, byte[] valueBytes, Cluster cluster) {
// 自定义逻辑
}
}3.3 ACK 与可靠性
acks | 行为 | 可靠性 | 性能 |
|---|---|---|---|
0 | 不等任何确认,发完即返回 | 最低 | 最高 |
1 | Leader 写入即返回,不等副本 | 中等 | 高 |
-1 (all) | 所有 ISR 副本都写入才返回 | 最高 | 较低 |
3.4 幂等与事务
幂等生产者(enable.idempotence=true):
每条消息带 (ProducerID, SequenceNumber)
Broker 去重: 同一个 ProducerID 的 SeqNum 只能递增事务(transactional.id):
java
producer.initTransactions();
producer.beginTransaction();
producer.send(...);
producer.send(...);
producer.commitTransaction(); // 或 abortTransaction()
// 支持跨分区跨 Topic 原子性 (EOS: Exactly Once Semantics)四、消费者
4.1 消费者组与 Rebalance
Consumer Group: app-log-processor
Partition: P0 P1 P2 P3 P4 P5
│ │ │ │ │ │
Consumer: C1 C1 C2 C2 C3 C3规则:
- 一个分区只能被组内一个消费者消费
- N 个分区最多 N 个消费者(多余的闲置)
- 消费者数量变化/超时 → 触发 Rebalance
4.2 Rebalance 协议演进
| 版本 | 策略 | 特点 |
|---|---|---|
| Eager (旧) | Stop-The-World | 所有消费者停止,重新分配,短暂不可用 |
| Cooperative (2.4+) | 增量分配 | 只迁移必要的分区,不停止所有消费者 |
4.3 Offset 管理
__consumer_offsets (内部 Topic, 默认 50 分区)
存储: (GroupID, Topic, Partition) → Offset
自动提交 (enable.auto.commit=true):
按 auto.commit.interval.ms (默认 5s)
手动提交:
consumer.commitSync() ← 同步, 阻塞, 重试
consumer.commitAsync() ← 异步, 不阻塞, 不重试4.4 消费语义
| 语义 | 实现方式 | 适用场景 |
|---|---|---|
| At Most Once | 先提交 offset 再处理 | 可丢数据(监控日志) |
| At Least Once (默认) | 先处理再提交 offset | 可能重复(需幂等) |
| Exactly Once | 事务 + 幂等 | 严格一致性 |
五、副本与 ISR
5.1 ISR 机制
Partition P0:
Leader (Broker 1) ← Producer 只写 Leader
Replica (Broker 2) ← ISR (In-Sync Replica)
Replica (Broker 3) ← ISR
ISR: 与 Leader 保持同步的副本集合
replica.lag.time.max.ms (默认 30s) 内追上即保持判断依据:不是按消息条数落后,而是按时间。只要副本在 replica.lag.time.max.ms 内追上了 Leader,就在 ISR 中。这避免了瞬时写入高峰导致副本被频繁踢出。
5.2 Leader 选举
优先副本选举 (Preferred Replica):
每个分区的第一个副本是 "Preferred Leader"
auto.leader.rebalance.enable=true 时自动均衡
Unclean Leader 选举:
unclean.leader.election.enable=false (默认)
ISR 为空时 → 宁可不可用也不选落后副本
unclean.leader.election.enable=true → 可用性优先,可能丢数据5.3 LEO vs HW
Leader: [m1][m2][m3][m4][m5] LEO=5
Follower A: [m1][m2][m3] LEO=3
Follower B: [m1][m2][m3][m4] LEO=4
HW (High Watermark) = min(LEO of all ISR) = 3
→ 消费者最多能读到 m3 (HW 之前的消息已全部副本确认)六、Consumer Group Rebalance
6.1 Rebalance 触发条件
mermaid
sequenceDiagram
participant C1 as Consumer 1
participant C2 as Consumer 2
participant G as GroupCoordinator<br/>(某Broker)
participant T as Topic Partitions
Note over C1,C2: 正常消费中
C1--xC1: 崩溃!
Note over G: heartbeat 超时<br/>session.timeout.ms
G->>G: 触发 Rebalance
G->>C2: 撤销当前分区分配
G->>C2: 重新分配: P0,P1,P2 全给 C2触发条件:
- Consumer 加入或离开 Group
max.poll.interval.ms超时(两次 poll 间隔过长)session.timeout.ms心跳超时(Consumer 假死/真死)- Topic 分区数变化
6.2 Eager vs Cooperative Rebalance
| 特性 | Eager (老) | Cooperative (新, 2.4+) |
|---|---|---|
| 行为 | Stop-The-World:所有 Consumer 先释放全部分区 | 增量分配:只撤销需要重新分配的分区 |
| 影响 | 短暂消费中断 | 零停机,未变分区继续消费 |
| 配置 | 默认 | partition.assignment.strategy=CooperativeStickyAssignor |
6.3 分区分配策略
| 策略 | 分配方式 | 适用 |
|---|---|---|
| RangeAssignor | 按 Topic 逐个分配,分区号连续分配给同一 Consumer | 默认,可能导致倾斜 |
| RoundRobinAssignor | 轮询分配所有分区 | 各 Consumer 分区数均匀 |
| StickyAssignor | 尽量保持上次分配结果 + 均衡 | Rebalance 时减少分区移动 |
| CooperativeStickyAssignor | Sticky + 增量 Rebalance | ✅ 生产推荐 |
七、Exactly-Once 语义
7.1 三层语义
mermaid
flowchart LR
subgraph Semantics["消息传递语义"]
A["At Most Once<br/>最多一次<br/>可能丢消息"]
B["At Least Once<br/>至少一次<br/>可能重复"]
C["Exactly Once<br/>精确一次<br/>不丢不重"]
end
A -->|"acks=0, 不重试"| Producer
B -->|"acks=all, 重试"| Producer
C -->|"幂等+事务"| Producer7.2 幂等 Producer(Idempotent)
java
// 开启幂等:Producer 自动去重
props.put("enable.idempotence", true);
// 内部机制:
// - 每个 Producer 分配 PID (Producer ID)
// - 每条消息带 sequence number(递增)
// - Broker 检查 (PID, TopicPartition, SeqNo) 是否重复
// - 自动设置 acks=all, max.in.flight=5, retries=MAX_INT7.3 事务(Transactions)
跨分区原子写入 + consume-transform-produce 精确一次:
java
// 初始化事务 Producer
props.put("transactional.id", "tx-app-1");
producer.initTransactions();
// 事务性写入
producer.beginTransaction();
try {
producer.send(new ProducerRecord<>("topic-A", key, valueA));
producer.send(new ProducerRecord<>("topic-B", key, valueB));
producer.commitTransaction();
} catch (Exception e) {
producer.abortTransaction();
}
// Consume-Transform-Produce 模式:
// consumer 读取 → 处理 → producer 事务写入 → 提交消费 offset
producer.sendOffsetsToTransaction(offsets, consumerGroup);八、性能参数调优
6.1 Broker 端
| 参数 | 说明 | 建议 |
|---|---|---|
num.network.threads | 网络处理线程 | CPU 核数 |
num.io.threads | 磁盘 I/O 线程 | CPU 核数 × 2 |
num.partitions | 总分区数 | 单 Broker ≤ 4000 |
log.segment.bytes | 日志段大小 | 1GB |
log.retention.hours | 保留时间 | 按需求 (168=7天) |
socket.send.buffer.bytes | 发送缓冲区 | 1MB |
socket.receive.buffer.bytes | 接收缓冲区 | 1MB |
6.2 Producer 端
| 参数 | 说明 | 建议 |
|---|---|---|
batch.size | 一批消息大小 | 16KB-1MB |
linger.ms | 等待攒批时间 | 5-100ms |
buffer.memory | 总缓冲区 | 32MB-128MB |
compression.type | 压缩算法 | lz4 或 snappy |
max.in.flight.requests.per.connection | 最大在途请求 | ≤ 5(开启幂等时 ≤ 5) |
6.3 Consumer 端
| 参数 | 说明 | 建议 |
|---|---|---|
fetch.min.bytes | 最少拉取字节 | 1MB |
fetch.max.wait.ms | 等待攒够 fetch.min.bytes 的最长时间 | 500ms |
max.partition.fetch.bytes | 每分区最大拉取 | 10MB |
max.poll.records | 单次 poll 最多条数 | 500-2000 |
max.poll.interval.ms | 两次 poll 最大间隔(超时触发 rebalance) | 5min |
九、常见问题
7.1 消息丢失场景
| 场景 | 原因 | 解决方案 |
|---|---|---|
| Producer | acks=0 不等待 | acks=-1 + min.insync.replicas=2 |
| Broker | Leader 宕机 + unclean 选举 | unclean.leader.election.enable=false |
| Consumer | 先提交 offset 后处理 | 先处理再提交(At Least Once) |
7.2 消息重复与幂等
重复消费不可避免 (At Least Once), 需业务幂等:
1. 唯一键去重: 每条消息带 UUID/业务 ID, DB 唯一约束
2. 版本号: Redis SETNX / DB version 字段
3. 事务标记: 处理前记录 (txID, 状态), 避免重复处理7.3 消息积压
| 原因 | 处理 |
|---|---|
| 消费速度跟不上 | 增加消费者(不超过分区数) |
| 消息倾斜 | 调整分区策略使 key 更均匀 |
| 下游慢/异常 | 熔断降级,临时扩容 |
| 消费代码问题 | 排查耗时逻辑,优化处理速度 |
7.4 消息顺序保证
Kafka 只保证单个分区内消息有序:
需要全局有序 → 1 个分区 (牺牲并行度)
需要按用户有序 → 按 user_id 分区十、Kafka vs 其他 MQ
| 维度 | Kafka | RabbitMQ | RocketMQ | Pulsar |
|---|---|---|---|---|
| 吞吐量 | 极高 (百万/秒) | 中等 (万/秒) | 高 (十万/秒) | 极高 |
| 延迟 | 毫秒级 | 微秒级 | 毫秒级 | 毫秒级 |
| 存储 | 磁盘持久化, 长期 | 内存+磁盘, 消费即删 | 磁盘持久化 | 分层存储 |
| 消费模型 | 拉取 (Pull) | 推送 (Push) | 拉取 (Pull) | 拉取 (Pull) |
| 可靠性 | 高 (副本+ISR) | 极高 | 高 (同步刷盘) | 高 |
| 顺序 | 分区内有序 | 队列内有序 | 队列内有序 | 分区内有序 |
| 事务 | ✅ (0.11+) | ❌ | ✅ | ❌ |
| 适用场景 | 日志、流、大数据管道 | 业务消息、RPC | 电商、金融 | 云原生、多租户 |
| 多租户 | ❌ (需额外隔离) | 部分(vhost) | ❌ | ✅ 原生 |
九、Kafka 生态
| 组件 | 作用 |
|---|---|
| Kafka Connect | 数据导入/导出连接器(DB → Kafka → ES/HDFS) |
| Kafka Streams | 轻量级流处理库(Java,无需独立集群) |
| KSQL (ksqlDB) | SQL 流处理(SELECT ... FROM stream WHERE ...) |
| Schema Registry | Avro/Protobuf Schema 管理 |
| MirrorMaker 2 | 跨集群数据复制 |
| KRaft (3.3+ 生产可用) | 去 ZooKeeper,元数据由 Kafka 自身 Raft 协议管理 |
十、Kafka vs RocketMQ vs Pulsar — 三大消息队列对比
三家都是"分布式消息队列",但设计哲学大不相同。选型不是在选"谁更好",而是在选"谁更适合你的场景":
| 维度 | Kafka | RocketMQ | Pulsar |
|---|---|---|---|
| 底层存储 | 日志段文件 (Segment) | CommitLog + ConsumeQueue | BookKeeper (Ledger) |
| 消息模型 | 分区 (Partition) | 队列 (MessageQueue) | Topic → Segment |
| 消费模型 | 拉 (Pull) 为主 | 拉+推 (Pull+Push) | 拉 (Pull) |
| 延迟消息 | ❌ 原生不支持 | ✅ 18 个延迟等级 | ✅ 任意延迟 |
| 事务消息 | ✅ Exactly-Once (幂等+事务) | ✅ 半消息回查机制 | ✅ |
| 顺序消息 | 🟡 分区内有序 | ✅ 局部有序 | ✅ |
| 流量控制 | 仅消费者 | 生产者+消费者 | 生产者+消费者 |
| 多租户 | ❌ 弱 | 🟡 通过命名空间 | ✅ 原生 Namespace |
| 存储与计算分离 | ❌ 绑定 | ❌ 绑定 | ✅ Broker 无状态 |
| 协议兼容 | 自有协议 | 自有协议 | Kafka/Pulsar/AMQP/RabbitMQ |
| 运维复杂度 | 🟡 中等 | 🟡 中等 | ❌ 高(多组件) |
| 社区/生态 | ✅✅ 最大 | ✅ 国内强 | 🟡 追赶中 |
| 典型用户 | LinkedIn, Uber, Netflix | 阿里, 滴滴, 美团 | Yahoo, Verizon, Tencent |
ISR 机制与数据可靠性
ISR(In-Sync Replicas)是 Kafka 保证数据可靠性的核心机制,但它的行为受到多个参数的影响:
mermaid
flowchart TB
Leader["Partition Leader<br/>接收所有写入"]
subgraph ISR["ISR (In-Sync Replicas)"]
F1["Follower 1<br/>拉取延迟 < replica.lag.time.max.ms(30s)"]
F2["Follower 2<br/>拉取延迟 < 30s"]
end
subgraph OSR["OSR (Out-of-Sync Replicas)"]
F3["Follower 3<br/>拉取延迟 > 30s ❌ 被踢出 ISR"]
end
Leader --> F1
Leader --> F2
Leader -.->|"延迟过大→踢出"| F3可靠性配置速查:
场景: 绝对不能丢消息
acks = -1 (all)
min.insync.replicas = 2 (至少 2 个 ISR 确认)
unclean.leader.election.enable = false (非 ISR 副本不能选为 Leader)
场景: 追求吞吐,偶尔丢可忍
acks = 1 (Leader 确认即可)
compression.type = lz4 (压缩减轻网络 I/O)
场景: 极低延迟 (如实时推荐)
acks = 0 (不等确认)
linger.ms = 0 (立即发送,不等批次)ISR 的核心权衡:
min.insync.replicas越大越可靠,但可用性越低——如果 ISR 数量降到此值以下,写入会直接拒绝。一般生产环境设replication.factor=3, min.insync.replicas=2,这样可以容忍 1 个节点故障。
登录后即可发表评论 👇