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)是为了批量压缩 + 零拷贝读放大,本质是为追求高吞吐
结论 B
Key 是消息结构里最关键的字段之一,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 = 这个快递在仓库货架上的位置编号——同一分区内的位置唯一
STEP 3 · 底层
两代消息格式对比
| 维度 | MessageSet(0.10-) | RecordBatch(0.11+) |
|---|---|---|
| 压缩粒度 | 每条独立 | 整批 |
| 位置表示 | 每条 offset | baseOffset + offsetDelta |
| 校验 | 单条 CRC(可选) | 整批 CRC(强制) |
| 元数据 | 少 | 支持 Headers |
| 版本 | Magic 0/1 | Magic 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) % numPartitionsQ: 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 张卡,点击翻面