Kafka 核心

分布式消息队列的日志内核——不可变追加日志、分区、消费者组偏移管理、ISR 复制。

概念

Kafka 是一个分布式流处理平台,核心是一个”不可变、顺序追加的分布式日志”。它的设计哲学极为精简:所有数据按顺序追加到分区日志文件中,每个消息由 (offset) 唯一索引。消费者通过记录已消费的偏移量 (offset) 实现消费进度管理。Kafka 用 Java 编写,部分组件 (如存储层) 用 C++ 重写 (Tiered Storage)。

核心组件

组件职责关键机制
Topic消息的逻辑分类可分多个分区
Partition水平扩展的基本单位顺序追加日志, 不可变
Segment分区文件的物理分段每个 segment 是一个文件
Offset消息在分区内的唯一 ID单调递增 (per partition)
Consumer Group分布式消费组自动负载均衡
ISRIn-Sync Replicas消息不丢失的复制保证
Leader / Follower每个分区的 leader 处理读写Controller 管理 leader 选举

分区日志模型

Topic "user-events" (3 partitions):

Partition 0: [msg0:offset=0] [msg1:offset=1] [msg2:offset=2] [msg3:offset=3] [msg4:offset=4] ...
Partition 1: [msg0:offset=0] [msg1:offset=1] [msg2:offset=2] ...
Partition 2: [msg0:offset=0] [msg1:offset=1] [msg2:offset=2] [msg3:offset=3] ...

每个 Partition 是物理上一个目录:
 /data/user-events-0/
 00000000000000000000.log (segment: offset 0 → 1000)
 00000000000000000000.index (稀疏索引)
 00000000000000001000.log (segment: offset 1000 → 2000)
 00000000000000001000.index

消费者组消费:
 Group "analytics":
 Consumer A → Partition 0 (offset=2, 消费中)
 Consumer B → Partition 1 (offset=5, 消费中)
 Consumer C → Partition 2 (offset=0, 消费中)

 偏移量存储在 __consumer_offsets topic 中 (Kafka 内部主题)

Consumer Group 偏移管理

消费者组协调:
 消费者组通过 __consumer_offsets topic 持久化每个分区的消费偏移量

 消息投递语义:
 At-most-once:
 fetch 消息 → commit offset → 处理消息
 (处理失败则丢消息)

 At-least-once:
 fetch 消息 → 处理消息 → commit offset
 (处理失败不 commit, 下次重收; 可能重复处理)

 Exactly-once:
 Kafka 0.11+ 支持事务:
 beginTransaction()
 producer.send("A"), producer.send("B")
 consumer.commit() // 原子提交: 生产和消费的位移一起提交
 commitTransaction()

 消费者 Rebalance:
 新消费者加入或退出 → 分区重新分配
 Rebalance 期间消费暂停
 常见策略: Range, RoundRobin, Sticky

ISR 复制机制

假设 Partition 0 有 3 个副本:

 Leader (Broker 1): [0][1][2][3][4][5][6][7][8][9] (LEO=10, HW=8)
 Follower A (Broker 2): [0][1][2][3][4][5][6][7][8] (LEO=9)
 Follower B (Broker 3): [0][1][2][3][4][5][6][7][8] (LEO=9)

 LEO (Log End Offset): 该副本最后一条消息的 offset + 1
 HW (High Watermark): 所有 ISR 副本中最小的 LEO

 消费者只能读到 HW 之前的数据 (offset < HW)
 // 保证消费者不会读到可能未确认的数据

 ISR (In-Sync Replicas):
 与 Leader 保持同步 (滞后不超过 replica.lag.time.max.ms) 的副本集合
 ISR = [Broker1, Broker2, Broker3] (所有副本都同步)

 如果 Follower B 滞后:
 ISR = [Broker1, Broker2]
 HW 变为 9 (Broker2 的 LEO)
 acks=all 时 producer 只等 ISR 中的 follower 确认