跳到正文
SL Blog 技术探索 · 工程实践 · AI 时代思考
返回
📨 Kafka 3.9 教程 · 5 / 10 查看系列简介 →

Kafka 3.9 教程(四):实现原理

文章目录
  1. 目录
  2. 4.1 网络层
  3. 4.1.1 通信协议
  4. 4.1.2 请求处理流程
  5. 4.1.3 零拷贝技术(Zero-Copy)
  6. 4.1.4 请求类型
  7. 4.2 消息格式
  8. 4.2.1 消息批次(Record Batch)
  9. 4.2.2 消息头部(Headers)
  10. 4.2.3 消息压缩
  11. 4.2.4 消息格式版本
  12. 4.3 日志存储
  13. 4.3.1 日志目录结构
  14. 4.3.2 日志段(Log Segment)
  15. 4.3.3 日志段滚动
  16. 4.3.4 索引机制
  17. 4.3.5 日志清理策略
  18. 4.4 数据分布与 Leader 选举
  19. 4.4.1 分区分配
  20. 4.4.2 Controller
  21. 4.4.3 Leader 选举流程
  22. 4.4.4 ISR(同步副本)管理
  23. 4.4.5 数据一致性保证
  24. 4.4.6 故障恢复
  25. 附录:关键指标监控
  26. Broker 指标
  27. 分区指标
  28. 下一步学习

📌 本文是「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_ISRLeader 选举
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最低
gzipGZIP 压缩2-3x中等
snappyGoogle Snappy2x低
lz4LZ4 压缩2x最低
zstdFacebook Zstandard3-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.bytes1GB
时间超限log.roll.ms7天
索引大小超限log.index.size.max.bytes10MB

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   # 无大小限制

删除流程:

  1. 检查日志段是否过期
  2. 标记可删除的日志段
  3. 删除文件或标记为删除

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 职责:

  1. 选举分区 Leader
  2. 跟踪 Broker 上下线
  3. 更新分区副本分配
  4. 管理主题创建/删除

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 保持同步的副本集合。

判断同步条件:

  1. 时间条件:replica.lag.time.max.ms 内有 fetch 响应
  2. 数量条件:落后 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不等待发送即成功
1Leader 确认Leader 写入成功
allISR 全部确认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
ISRExpansionRateISR 扩展速率持续 > 0
ISRShrinkRateISR 收缩速率持续 > 0
LeaderCountLeader 分区总数异常波动
TotalTimeMs请求处理总时间> 100ms

分区指标

指标说明
LogEndOffset分区末尾偏移量
LogStartOffset分区起始偏移量
NumLogSegments日志段数量
Size分区数据大小

下一步学习


RAG 智能问答

针对本文继续提问:《Kafka 3.9 教程(四):实现原理》