---
title: Kafka 消息结构 + Offset
type: concept
domain: mq
tags: [kafka, message-structure, recordbatch, offset, messagekey, headers]
status: learning
created: 2026-09-22
last_reviewed: 2026-09-22
method: feynman
related_questions:
  - "Kafka 的消息的结构是什么样的"
  - "Kafka 中的 Offset 是什么？"
related_knowledge:
  - ./kafka-partition-performance.md
  - ./kafka-reliability.md
visualizations: []
anki_cards: 5
interview_rounds:
  - "2026-09-22-round-2-Q1"
  - "2026-09-22-round-2-Q2"
---

# Kafka 消息结构 + Offset

> 消息结构决定"消息是什么"，Offset 决定"消息在哪"——两者共同构成 Kafka 的消息寻址模型。

## 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 已提交"

## 关键结论

- **结论 A**：Kafka 消息结构的两次演进（MessageSet → RecordBatch）是为了**批量压缩 + 零拷贝读放大**，本质是为追求高吞吐
- **结论 B**：`Key` 是消息结构里**最关键的字段之一**，`hash(key) % numPartitions` 决定消息路由，是"同一业务实体顺序消费"的基石
- **结论 C**：Offset 只在**单 Partition 内**唯一，Kafka 没有全局唯一的 Message ID
- **结论 D**：Offset 提交是 **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+） |
|------|-------------------|---------------------|
| 压缩粒度 | 每条独立 | 整批 |
| 位置表示 | 每条 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/<group>/offsets/<topic>/<partition>`
- 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` 里
- **提交**：手动 + 幂等才是王道

## 常见误区

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

## 延伸追问

1. **Kafka 0.11 之后为什么把 MessageSet 改成 RecordBatch？**
   - 批量压缩减少 CPU 开销；baseOffset 让单条消息更紧凑；支持按批读放大（零拷贝）。
2. **Key 是 null 和 Key="" 有什么区别？**
   - `null` 走**粘性分区**（sticky partitioner，轮流选）；`""` 走 hash 路由到固定分区。
3. **消息体超过 `max.message.bytes` 会怎样？**
   - Producer 端直接抛异常，不会截断；需要调 producer 和 broker 两端配置。
4. **Offset 会被复用吗？Compaction 删除的消息怎么办？**
   - 不会复用；Compaction 保留每个 key 最新值，其他物理删除，新消息写入不复用旧 offset。
5. **`__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，仅首次生效
```

## 关联题目（题库）

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

## 关联知识

- [Kafka 分区与性能](./kafka-partition-performance.md)
- [Kafka 消息可靠性](./kafka-reliability.md)
- [主题地图](./_moc.md)

## Anki 候选卡片

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

---

*最后更新：2026-09-22*
