📌 本文是「Kafka 3.9 教程」系列第 四 篇 · 系列目录
本文档基于 Apache Kafka 3.9 官方文档整理
目录
4.1 网络层
4.1.1 通信协议
Kafka 使用高性能 TCP 网络协议进行客户端与 Broker 之间的通信。
协议特点:
- 二进制协议:紧凑、高效
- 请求/响应模式:同步和异步操作
- 批量处理:多个请求可以合并发送
- 压缩支持:支持 gzip、snappy、lz4、zstd
4.1.2 请求处理流程
客户端请求
│
▼
┌─────────────────────────────────────────────┐
│ 1. TCP 连接建立 │
│ - 三次握手 │
│ - SSL/TLS 握手(如启用) │
└─────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────┐
│ 2. 请求发送 │
│ - 序列化请求 │
│ - 压缩(如启用) │
│ - 发送到 Socket │
└─────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────┐
│ 3. Broker 处理 │
│ - 反序列化请求 │
│ - 权限校验 │
│ - 请求队列排队 │
│ - 处理请求 │
│ - 响应生成 │
└─────────────────────────────────────────────┘
│
▼
┌─────────────────────────────────────────────┐
│ 4. 响应返回 │
│ - 序列化响应 │
│ - 发送回客户端 │
└─────────────────────────────────────────────┘
4.1.3 零拷贝技术(Zero-Copy)
Kafka 使用零拷贝技术减少数据传输开销。
传统数据传输:
磁盘 ──▶ 内核缓冲区 ──▶ 用户空间 ──▶ Socket 缓冲区 ──▶ 网卡
│ │ │
4次拷贝 2次上下文切换
零拷贝技术:
磁盘 ──▶ 内核缓冲区 ──▶ 网卡
│ │
2次拷贝 直接传输(sendfile)
│
1次上下文切换
带来的好处:
- 减少 CPU 开销
- 减少内存拷贝
- 提高吞吐量
- 降低延迟
4.1.4 请求类型
| 请求类型 | 说明 |
|---|---|
| PRODUCE | 生产消息 |
| FETCH | 消费消息 |
| LIST_OFFSETS | 查询偏移量 |
| METADATA | 查询集群元数据 |
| LEADER_AND_ISR | Leader 选举 |
| STOP_REPLICA | 停止副本 |
| CONTROLLED_SHUTDOWN | 优雅关闭 |
4.2 消息格式
4.2.1 消息批次(Record Batch)
Kafka 0.11.0 引入了新的消息格式,支持幂等性和事务。
消息批次结构:
┌─────────────────────────────────────────────────────────┐
│ Record Batch │
├─────────────────────────────────────────────────────────┤
│ ┌─────────────────────────────────────────────────┐ │
│ │ 基础信息 │ │
│ │ - baseOffset: int64 │ │
│ │ - batchLength: int32 │ │
│ │ - partitionLeaderEpoch: int32 │ │
│ │ - magic: int8 (消息格式版本) │ │
│ │ - crc: int32 (校验和) │ │
│ │ - attributes: int16 (压缩、时间戳类型等) │ │
│ │ - lastOffsetDelta: int32 │ │
│ │ - firstTimestamp: int64 │ │
│ │ - maxTimestamp: int64 │ │
│ │ - producerId: int64 │ │
│ │ - producerEpoch: int16 │ │
│ │ - sequenceNumber: int32 │ │
│ └─────────────────────────────────────────────────┘ │
│ ┌─────────────────────────────────────────────────┐ │
│ │ 消息列表 │ │
│ │ - length: int32 │ │
│ │ - attributes: int8 │ │
│ │ - timestampDelta: int64 │ │
│ │ - offsetDelta: int32 │ │
│ │ - key: bytes │ │
│ │ - value: bytes │ │
│ │ - headers: [key, value] 列表 │ │
│ └─────────────────────────────────────────────────┘ │
└─────────────────────────────────────────────────────────┘
4.2.2 消息头部(Headers)
Kafka 0.11.0 引入了消息头部,用于存储元数据。
// 生产者设置头部
ProducerRecord<String, String> record = new ProducerRecord<>(
"my-topic",
"key",
"value"
);
// 添加头部
record.headers().add("trace-id", "12345".getBytes());
record.headers().add("content-type", "application/json".getBytes());
// 消费者读取头部
ConsumerRecord<String, String> record = ...;
Iterable<Header> headers = record.headers();
for (Header header : headers) {
System.out.println(header.key() + ": " + new String(header.value()));
}
4.2.3 消息压缩
Kafka 支持多种压缩算法:
| 压缩类型 | 说明 | 压缩率 | CPU 开销 |
|---|---|---|---|
| none | 不压缩 | 1x | 最低 |
| gzip | GZIP 压缩 | 2-3x | 中等 |
| snappy | Google Snappy | 2x | 低 |
| lz4 | LZ4 压缩 | 2x | 最低 |
| zstd | Facebook Zstandard | 3-4x | 中等 |
压缩位置:
- 生产者压缩:在发送前压缩整个批次
- Broker 压缩:可以配置 Broker 重新压缩
最佳实践:
- 高吞吐量场景:使用 lz4
- 高压缩率场景:使用 zstd
- 兼容性优先:使用 gzip
4.2.4 消息格式版本
| 版本 | 说明 | 引入版本 |
|---|---|---|
| 0 | 原始格式 | 0.8.x |
| 1 | 增加时间戳 | 0.10.x |
| 2 | 支持幂等和事务 | 0.11.x |
| 3 | 保留字段扩展 | 2.x |
# 配置消息格式版本
log.message.format.version=2.8
4.3 日志存储
4.3.1 日志目录结构
/tmp/kafka-logs/
├── my-topic-0/
│ ├── 00000000000000000000.log
│ ├── 00000000000000000000.index
│ ├── 00000000000000000000.timeindex
│ ├── 00000000000000100000.log
│ ├── 00000000000000100000.index
│ └── 00000000000000100000.timeindex
├── my-topic-1/
│ └── ...
└── my-topic-2/
└── ...
4.3.2 日志段(Log Segment)
每个分区对应一个目录,目录中包含多个日志段。
日志段组成:
┌─────────────────────────────────────────┐
│ Log Segment │
├─────────────────────────────────────────┤
│ .log 文件 │
│ - 存储实际消息内容 │
│ - 格式:消息批次 │
├─────────────────────────────────────────┤
│ .index 文件 │
│ - 偏移量索引 │
│ - 格式:<offset, position> │
│ - 稀疏索引 │
├─────────────────────────────────────────┤
│ .timeindex 文件 │
│ - 时间戳索引 │
│ - 格式:<timestamp, offset> │
│ - 支持时间查询 │
├─────────────────────────────────────────┤
│ .txnindex 文件(2.4+) │
│ - 事务索引 │
└─────────────────────────────────────────┘
4.3.3 日志段滚动
日志段在以下条件触发滚动:
| 条件 | 配置项 | 默认值 |
|---|---|---|
| 日志段大小超限 | log.segment.bytes | 1GB |
| 时间超限 | log.roll.ms | 7天 |
| 索引大小超限 | log.index.size.max.bytes | 10MB |
4.3.4 索引机制
偏移量索引:
- 稀疏索引,每 4KB 创建一个索引条目
- 支持快速定位指定偏移量的消息
- 二分查找 + 顺序扫描
.index 文件内容示例:
┌──────────┬─────────┐
│ offset │ position│
├──────────┼─────────┤
│ 0 │ 0 │
│ 100 │ 500 │
│ 200 │ 1200 │
│ 300 │ 2100 │
└──────────┴─────────┘
时间索引:
- 按时间戳建立索引
- 支持时间范围查询
- 查找指定时间范围内的消息
4.3.5 日志清理策略
Delete(删除策略)
# 配置
cleanup.policy=delete
log.retention.hours=168 # 7天
log.retention.bytes=-1 # 无大小限制
删除流程:
- 检查日志段是否过期
- 标记可删除的日志段
- 删除文件或标记为删除
Compact(压缩策略)
# 配置
cleanup.policy=compact
log.compaction.records.retention=100
min.compaction.lag.ms=0
压缩原理:
- 只保留每个 key 的最新值
- 删除 tombstone 标记的消息
- 保持消息顺序
压缩前: 压缩后:
key: A, value: 1 key: A, value: 4
key: B, value: 2 key: B, value: 3
key: A, value: 3 key: C, value: 5
key: B, value: 3
key: C, value: 5
key: A, value: 4
应用场景:
- 变更数据捕获(CDC)
- 配置更新
- 用户状态追踪
4.4 数据分布与 Leader 选举
4.4.1 分区分配
Broker 间分配:
- 使用分区分配策略
- 考虑副本分散和负载均衡
// 分区分配示例
// 假设有 4 个 Broker (0-3),主题有 9 个分区,副本因子 3
// 分配结果:
// Partition 0: [0, 1, 2]
// Partition 1: [1, 2, 3]
// Partition 2: [2, 3, 0]
// Partition 3: [3, 0, 1]
// ...
4.4.2 Controller
Controller 是 Kafka 集群的核心组件,负责管理分区 Leader 选举和元数据。
Controller 职责:
- 选举分区 Leader
- 跟踪 Broker 上下线
- 更新分区副本分配
- 管理主题创建/删除
KRaft 模式:
- Controller 角色由 Kafka 集群自身担任
- 使用 Raft 协议选举 Controller
- 不再依赖 ZooKeeper
4.4.3 Leader 选举流程
Leader 故障
│
▼
┌───────────────────────┐
│ 检测到 Leader 故障 │
│ (通过 ZooKeeper │
│ 或 Controller) │
└───────────────────────┘
│
▼
┌───────────────────────┐
│ 从 ISR 中选举 │
│ (In-Sync Replicas) │
└───────────────────────┘
│
▼
┌───────────────────────┐
│ 更新 ZooKeeper/ │
│ KRaft 元数据 │
└───────────────────────┘
│
▼
┌───────────────────────┐
│ 通知所有 Broker │
│ 更新 Leader 信息 │
└───────────────────────┘
│
▼
┌───────────────────────┐
│ 客户端重定向 │
│ (返回新 Leader) │
└───────────────────────┘
4.4.4 ISR(同步副本)管理
ISR 定义: 与 Leader 保持同步的副本集合。
判断同步条件:
- 时间条件:
replica.lag.time.max.ms内有 fetch 响应 - 数量条件:落后 Leader 的消息数在阈值内
ISR 状态示例:
Partition 0 的 Leader 在 Broker 0
ISR = [Broker 0, Broker 1] // Broker 2 落后太多,被移出 ISR
正常情况:
┌─────────────────────────────────────────────────────┐
│ Broker 0 (Leader) ────────────────────┐ │
│ │ │ │
│ │ 同步 │ 同步 │
│ ▼ ▼ │
│ Broker 1 (Follower, ISR) Broker 2 (Follower) │
│ 落后 → 移出 ISR │
└─────────────────────────────────────────────────────┘
4.4.5 数据一致性保证
写一致性级别:
| acks | 含义 | 写入时机 |
|---|---|---|
| 0 | 不等待 | 发送即成功 |
| 1 | Leader 确认 | Leader 写入成功 |
| all | ISR 全部确认 | ISR 全部写入成功 |
读一致性:
- Read-your-writes:写入后能立即读到
- Monotonic reads:不会读到旧数据
4.4.6 故障恢复
场景1:Follower 故障
故障:Broker 2 上的 Follower 崩溃
│
▼
检测:Leader 超过 replica.lag.time.max.ms 未收到心跳
│
▼
移除:从 ISR 中移除 Broker 2
│
▼
恢复:Broker 2 重启后追赶 Leader 数据
│
▼
加入:Broker 2 同步完成后加入 ISR
场景2:Leader 故障
故障:Broker 0 上的 Leader 崩溃
│
▼
检测:Controller 检测到 Leader 失联
│
▼
选举:从 ISR 中选择新 Leader(Broker 1)
│
▼
切换:元数据更新,客户端重定向
│
▼
恢复:旧 Leader(Broker 0)恢复后作为 Follower 同步
附录:关键指标监控
Broker 指标
| 指标 | 说明 | 告警阈值 |
|---|---|---|
| UnderReplicatedPartitions | 副本不足的分区数 | > 0 |
| ISRExpansionRate | ISR 扩展速率 | 持续 > 0 |
| ISRShrinkRate | ISR 收缩速率 | 持续 > 0 |
| LeaderCount | Leader 分区总数 | 异常波动 |
| TotalTimeMs | 请求处理总时间 | > 100ms |
分区指标
| 指标 | 说明 |
|---|---|
| LogEndOffset | 分区末尾偏移量 |
| LogStartOffset | 分区起始偏移量 |
| NumLogSegments | 日志段数量 |
| Size | 分区数据大小 |
下一步学习
- 第五章:运维指南 - 集群运维最佳实践
- 第六章:安全认证 - 安全配置详解
- 第七章:Kafka Connect - 数据集成框架
- 第八章:Kafka Streams - 流处理库