Kafka 消息结构 + Offset

mq 📚 learning kafka-message-structure · kafka · message-structure · recordbatch · offset · messagekey · headers

TL;DR(30 秒扫完)

  • 两代格式:0.10 之前 MessageSet(单条独立)→ 0.11+ RecordBatch(按批打包 + 批量压缩 + 整批 CRC)
  • 单条 Record 关键字段:Key(路由键,决定 partition)+ Value(消息体)+ Headers(元数据)+ Timestamp + offsetDelta
  • Offset 本质:Partition 内的单条消息索引,从 0 起单调递增,不可复用;(topic, partition, offset) 三元组唯一定位
  • Offset 存储演进:0.9 之前存 ZooKeeper → 0.9+ 存内置 Topic __consumer_offsets(50 分区)
  • 提交策略:enable.auto.commit=false + 业务完成后 commitSync 才能避免"处理失败但 offset 已提交"

关键结论

结论 AKafka 消息结构的两次演进(MessageSet → RecordBatch)是为了批量压缩 + 零拷贝读放大,本质是为追求高吞吐
结论 BKey 是消息结构里最关键的字段之一,hash(key) % numPartitions 决定消息路由,是"同一业务实体顺序消费"的基石
结论 COffset 只在单 Partition 内唯一,Kafka 没有全局唯一的 Message ID
结论 DOffset 提交是 Consumer Group 级别的行为,不是消费者实例自己维护

完整讲解(费曼四步)

STEP 1 · 概念

消息结构描述一条 Kafka 消息的字节布局;Offset 是消息在 Partition 内的唯一位置索引,从 0 起单调递增。

STEP 2 · 大白话

快递包裹比喻:
  • 单条 Record = 一个快递包裹,上面贴着运单号(Key + Timestamp + Headers)+ 内容(Value)+ 体积(length)
  • RecordBatch = 一整箱快递(打包发运),箱子上有总运单号区间(baseOffset + lastOffsetDelta)+ 校验码(CRC)+ 箱子编号(partitionLeaderEpoch)
  • Offset = 这个快递在仓库货架上的位置编号——同一分区内的位置唯一
Kafka 从"一件件发"升级到"整箱发",就是 0.11 那次演进。

STEP 3 · 底层

两代消息格式对比

维度MessageSet(0.10-)RecordBatch(0.11+)
压缩粒度每条独立整批
位置表示每条 offsetbaseOffset + offsetDelta
校验单条 CRC(可选)整批 CRC(强制)
元数据少支持 Headers
版本Magic 0/1Magic 2

单条 Record 字段(Magic 2)

attributes(压缩类型、isDeleted)
  + key.length + key            # 路由键
  + value.length + value        # 消息体
  + numHeaders + headers[]      # 用户自定义元数据(如 traceId)
  + timestamp + offsetDelta     # 相对 baseOffset 的偏移

RecordBatch 头部

baseOffset                # 本批起始偏移
batchLength
partitionLeaderEpoch      # 感知 Leader 切换
magicValue                # = 2
crc                       # 32 位,覆盖除自身外所有字节
attributes                # 压缩类型等
lastOffsetDelta / firstOffsetDelta
numRecords

Offset 的存储与提交

存储位置演进:
  • Kafka 0.9 之前:ZooKeeper /consumers//offsets//
  • Kafka 0.9+:内置 Topic __consumer_offsets,50 个分区,key = hash(group + topic + partition)
提交策略:
策略配置效果风险
自动提交enable.auto.commit=true(默认 5s)简单可能重复消费
commitSync手动调用强一致性能差
commitAsync手动调用性能好回调异常需处理
手动提交(推荐)enable.auto.commit=false + 业务后 commit最保险需业务幂等
auto.offset.reset:earliest / latest / none,仅在首次订阅且无 offset 记录时生效。

STEP 4 · 简化

一句话总结:Kafka 消息 = RecordBatch(打包)里的 Record(单条);Offset 是 Record 在 Partition 内的位置。 记忆口诀:
  • 消息结构:Key 定路由,Batch 走压缩,CRC 保完整
  • Offset:单区唯一,单调递增,Consumer Group 提交
  • 存储:0.9 之后都在 __consumer_offsets 里
  • 提交:手动 + 幂等才是王道

常见误区

说"Offset 是全局唯一的"
Offset 只在单 Partition 内唯一,全局定位要 (topic, partition, offset) 三元组
说"消息里有 Topic 字段"
Topic 属于 Partition 层组织概念,不在单条 Message 里
把 offset 说成"消息 ID"
Kafka 没有全局唯一 Message ID,只有 (topic, partition, offset)
说"消费者自己维护 offset"
Offset 提交是 Consumer Group 级别,不是消费者实例自己维护
认为 auto.offset.reset 决定所有消费起点
仅在首次订阅且无 offset 时生效

延伸追问

Kafka 0.11 之后为什么把 MessageSet 改成 RecordBatch?
批量压缩减少 CPU 开销;baseOffset 让单条消息更紧凑;支持按批读放大(零拷贝)。
Key 是 null 和 Key="" 有什么区别?
null 走粘性分区(sticky partitioner,轮流选);"" 走 hash 路由到固定分区。
消息体超过 max.message.bytes 会怎样?
Producer 端直接抛异常,不会截断;需要调 producer 和 broker 两端配置。
Offset 会被复用吗?Compaction 删除的消息怎么办?
不会复用;Compaction 保留每个 key 最新值,其他物理删除,新消息写入不复用旧 offset。
__consumer_offsets 为什么是 50 个分区?
平衡元数据管理与负载均衡;hash(group+topic+partition) % 50。

速查表

两代格式: MessageSet(0.10-) → RecordBatch(0.11+)
Record:   Key + Value + Headers + Timestamp + offsetDelta
Batch:    baseOffset + crc + attributes + numRecords
Offset:   单区唯一,单调递增,三元组唯一定位
存储:     0.9+ 存 __consumer_offsets(50 分区)
提交:     enable.auto.commit=false + 业务后 commitSync
reset:    earliest / latest / none,仅首次生效

Anki 候选卡片

Q: Kafka 0.11 之后消息格式叫什么?
A: RecordBatch(Magic 2),按批压缩 + 整批 CRC
Q: Kafka 消息结构里决定分区路由的字段?
A: Key,hash(key) % numPartitions
Q: Kafka 有没有全局唯一的 Message ID?
A: 没有,(topic, partition, offset) 三元组才唯一定位
Q: Kafka 0.9+ 消费者 Offset 存哪里?
A: 内置 Topic __consumer_offsets,50 个分区
Q: auto.offset.reset 什么时候生效?
A: 仅在消费者组首次订阅且无历史 offset 记录时生效

关联题目

  • ⚠️ 《Kafka 的消息的结构是什么样的》— Round 2 Q1, ⭐⭐(漏 Key,把 offset 说成唯一 ID)
  • ⚠️ 《Kafka 中的 Offset 是什么?》— Round 2 Q2, ⭐⭐⭐(漏 Consumer Group 提交和存储演进)

关联知识

消息 = RecordBatch 里的 Record;Offset 是 Partition 内的位置索引
快递包裹:Key 是运单号定路由,Batch 是整箱打包发运
✦ 记 忆 口 诀 ✦
Key 定路由,Batch 走压缩,Offset 单区唯一
关键可视化
RecordBatch + Record 结构
flowchart TD
  A[RecordBatch 头部] --> B[baseOffset + batchLength + magicValue]
  A --> C[crc 校验]
  A --> D[attributes 压缩类型]
  A --> E[numRecords]
  E --> F[Record 1]
  E --> G[Record 2]
  E --> H[Record N]
  F --> I[key 路由键]
  F --> J[value 消息体]
  F --> K[headers 元数据]
  F --> L[timestamp + offsetDelta]
Offset 存储演进
flowchart LR
  A[0.9 之前] --> B[ZooKeeper /consumers/offsets]
  C[0.9 之后] --> D[内置 Topic __consumer_offsets]
  D --> E[50 个分区]
  D --> F[key 是 group 加 topic 加 partition 的 hash]
知识关系

⬆️ 前置(Prerequisite)

暂无

🔄 延伸(Extends)

暂无

⚡ 对比(Contrast)

暂无
🎯 概念 📏 规则 ⚠️ 误区 🔍 追问 ✨ 口诀 共 0 张卡,点击翻面