Kafka 消息结构全解析:RecordBatch、Record 与 Headers 的二进制格式与源码实现
2026/9/10 13:28:35 网站建设 项目流程

Kafka 消息结构全解析:RecordBatch、Record 与 Headers 的二进制格式与源码实现

【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka

一条 Kafka 消息(Message / Record)在磁盘与网络上的真实形态是什么?本文以官方实现文档 docs/implementation/messages.md 与同目录下的 docs/implementation/message-format.md 为核心骨架,结合clients模块中org.apache.kafka.common.record包的真实源码,系统讲解 Kafka 消息的三段式结构、记录批(Record Batch)与单条记录(Record)的逐字段二进制格式、控制批(Control Batch)在事务与 KRaft 中的作用,以及旧版 message set 格式的来龙去脉。读完本文,你将具备直接阅读 Kafka 日志段二进制数据、理解生产/消费链路底层编码、排查消息格式兼容性问题的能力。

消息的三段式结构:头部 + 不透明 Key + 不透明 Value

官方实现文档messages.md对 Kafka 消息结构给出了最精炼的定义:一条消息由一个可变长度的头部(header)、一个可变长度的不透明 key 字节数组和一个可变长度的不透明 value 字节数组组成。头部字段的具体二进制布局在 docs/implementation/message-format.md 中描述——需要注意的是,Kafka 中"头部"分两层:记录批(RecordBatch)有自己的头部,批内的每条记录(Record)也有自己的头部,后者即我们常说的消息 Headers。

为什么 key 和 value 必须保持"不透明"?

文档给出的理由非常明确:

Leaving the key and value opaque is the right decision: there is a great deal of progress being made on serialization libraries right now, and any particular choice is unlikely to be right for all uses.

即:当前序列化库(Avro、Protobuf、JSON、自定义二进制等)发展迅速,没有任何一种序列化方案能适用于所有场景,因此 Kafka 核心把 key/value 一律视为字节数组,绝不感知其内部结构。序列化与反序列化完全由使用 Kafka 的应用程序自己负责——一个具体应用通常会约定某种序列化类型作为其使用规范的一部分。

这一设计在源码中体现得淋漓尽致:Record接口(clients/src/main/java/org/apache/kafka/common/record/internal/Record.java)中,key 与 value 的类型都是ByteBuffer,并仅通过keySize()/valueSize()(无 key 返回 -1)与key()/value()(可返回 null)暴露给调用方,不提供任何反序列化钩子。

// Record 接口(节选) long offset(); // 该记录在日志中的偏移 int sequence(); // 生产者分配的序列号(幂等/事务用) long timestamp(); // 记录时间戳 ByteBuffer key(); // 可为 null ByteBuffer value(); // 可为 null Header[] headers(); // magic < 2 时恒为空数组

RecordBatch:消息的批量容器与 NIO 读写入口

文档进一步说明:

TheRecordBatchinterface is simply an iterator over messages with specialized methods for bulk reading and writing to an NIOChannel.

RecordBatch本质上是一个"消息迭代器",并附带面向 NIOChannel的批量读写专用方法。对应实现为 clients/src/main/java/org/apache/kafka/common/record/internal/RecordBatch.java,它继承Iterable<Record>,同时暴露baseOffset()lastOffset()producerId()baseSequence()compressionType()isTransactional()isControlBatch()checksum()sizeInBytes()等批量级属性,以及writeTo(ByteBuffer)和延迟解压迭代器streamingIterator(BufferSupplier)

批的物理容器是 clients/src/main/java/org/apache/kafka/common/record/internal/MemoryRecords.java,它提供readableRecords(ByteBuffer)withRecords(...)withIdempotentRecords(...)withTransactionalRecords(...)withEndTransactionMarker(...)等工厂方法,覆盖普通、幂等、事务三种写入场景。

记录批(RecordBatch)的磁盘格式

message-format.md指出:消息(Records)总是以批(Batch)为单位写入,一个批包含一条或多条记录,极端情况下一个批也可以只含一条记录。批和记录各有自己的头部。RecordBatch 在磁盘上的逐字段布局如下(当前 magic 值为 2):

baseOffset: int64 batchLength: int32 partitionLeaderEpoch: int32 magic: int8 (current magic value is 2) crc: uint32 attributes: int16 bit 0~2: 0: no compression 1: gzip 2: snappy 3: lz4 4: zstd bit 3: timestampType bit 4: isTransactional (0 means not transactional) bit 5: isControlBatch (0 means not a control batch) bit 6: hasDeleteHorizonMs (0 means baseTimestamp is not set as the delete horizon for compaction) bit 7~15: unused lastOffsetDelta: int32 baseTimestamp: int64 maxTimestamp: int64 producerId: int64 producerEpoch: int16 baseSequence: int32 recordsCount: int32 records: [Record]

关键字段语义

  • baseOffset (int64):批内第一条记录的起始偏移。结合源码(RecordBatch.java)可知,对 magic 2+ 而言baseOffset()返回的是压缩前原始批的第一条偏移,即使压缩清掉部分记录也不会变化;而对 magic 0/1,获取 base offset 需要深遍历,因此文档建议优先使用更高效的lastOffset()
  • batchLength (int32):从本字段之后到批末尾的字节数。因此一个批在磁盘上的总大小 =batchLength + 12(12 = 8 字节 baseOffset + 4 字节 batchLength 自身)。
  • partitionLeaderEpoch (int32):分区 leader 的纪元号。
  • magic (int8):当前取值为 2,对应源码中的CURRENT_MAGIC_VALUE = MAGIC_VALUE_V2(另存在MAGIC_VALUE_V0 = 0MAGIC_VALUE_V1 = 1)。
  • crc (uint32)覆盖从 attributes 到批末尾的所有字节(即 CRC 之后的所有内容),使用 CRC-32C(Castagnoli)多项式。CRC 位于 magic 之后,因此客户端必须先解析 magic 字节,才能决定如何解释 batchLength 与 magic 之间的字节。partitionLeaderEpoch 字段不参与 CRC 计算——这是刻意为之:该字段由 broker 在接收每个批时赋值,若不排除它,每次赋值都需重算 CRC。

attributes 位域:压缩、时间戳类型、事务与控制批

attributes 的 16 位被细分为多个语义位:

含义取值说明
bit 0~2压缩算法0: 无压缩;1: gzip;2: snappy;3: lz4;4: zstd
bit 3timestampType区分 CreateTime 与 LogAppendTime
bit 4isTransactional1 表示该批属于某个事务
bit 5isControlBatch1 表示控制批
bit 6hasDeleteHorizonMs1 表示 baseTimestamp 被设置为压缩删除边界(delete horizon)
bit 7~15未使用保留位

压缩算法常量与编码对应关系在 clients/src/main/java/org/apache/kafka/common/record/internal/CompressionType.java 中定义。当启用压缩时,压缩后的记录数据被直接序列化在 recordsCount 字段之后——也就是说压缩针对的是整批记录的字节流,而不是逐条记录独立压缩。

与压缩相关的生产者侧实践

生产者端的compression.type配置(可取nonegzipsnappylz4zstd)最终会写入上述 attributes 的低 3 位。zstd 由 Kafka 2.1 起引入支持;压缩通常能显著降低网络与磁盘占用,但会增加生产端 CPU 开销与消费端解压开销,小消息批量场景下收益有限。

日志压缩(Log Compaction)对批的特殊处理

message-format.md用较大篇幅说明了压缩(此处指 topic 的 log compaction 清理,非上文压缩算法)与批格式的交互,理解这点对排查幂等/事务生产者在分区 leader 切换后的OutOfSequence错误至关重要:

  • 保留首尾 offset 与首尾 sequence:日志清理时,批的第一个和最后一个 offset/sequence 会被保留,因为日志重载时需要据此恢复生产者状态。若只保留首 sequence 而丢失末 sequence,分区 leader 故障后生产者可能遇到OutOfSequence错误;而 base sequence 必须保留用于重复消息检查——broker 校验 Produce 请求时,会比对请求中批的首尾 sequence 与该生产者上一次记录的 sequence 是否衔接。
  • 允许出现空批:当批内所有记录都被清理、但为了保留生产者最后的 sequence 时,日志中会残留空批(batch 无记录但依然存在)。
  • baseTimestamp 可能变化:压缩时 baseTimestamp 不保留,若批内首条记录被压缩掉,baseTimestamp 会随之改变。
  • delete horizon 机制:若批中包含 null 负载(tombstone)或中止事务标记的记录,压缩可能修改 baseTimestamp,将其设为这些记录"应被删除的时刻",并同时置位 attributes 的 bit 6(hasDeleteHorizonMs),从而告诉压缩器何时可以安全地清除这些记录。

对应的 API 是RecordBatch#deleteHorizonMs()(RecordBatch.java),返回OptionalLong,若批的首时间戳不是 delete horizon 则返回空。

控制批(Control Batch):事务与 KRaft 的底层信令

控制批(Control Batch)是 magic 2 引入的特殊批类型,其中只包含一条称为控制记录(control record)的记录。控制记录绝不会返回给应用程序,其用途有二:

  1. 消费者侧过滤:消费者借助控制记录过滤掉已中止(abort)的事务消息;
  2. KRaft 协议元数据:控制批被 KRaft 共识协议用于承载内部协议元数据。

控制记录的 key 遵循如下 schema:

version: int16 (current version is 0) type: int16 (the control record types are in the table below)

常规 topic(regular topics)当前定义的控制记录类型如下:

TypeNameDescription
0ABORT标记一个事务已中止
1COMMIT标记一个事务已提交

类型 0 和 1 用作事务消息协议的事务结束标记(end-of-transaction markers);类型 2 到 6 由 KRaft 共识协议内部使用。控制记录 value 的 schema 取决于类型,对客户端而言 value 是不透明的。

源码实现位于 clients/src/main/java/org/apache/kafka/common/record/internal/ControlRecordType.java:枚举定义了ABORT((short) 0)COMMIT((short) 1)以及UNKNOWN((short) -1)(用于表示客户端无法识别的控制类型,应被忽略),fromTypeId(short)负责类型 ID 到枚举的映射,遇到未注册的类型 ID 则返回UNKNOWN以保证向前兼容。

单条记录(Record)的二进制格式

批内的每条记录(magic 2+)遵循如下紧凑布局,核心设计思路是用相对量(delta)替代绝对量,配合 varint 压缩整数编码,最大化节省空间:

length: varint attributes: int8 bit 0~7: unused timestampDelta: varlong offsetDelta: varint keyLength: varint key: byte[] valueLength: varint value: byte[] headersCount: varint Headers => [Header]

各字段语义:

  • length (varint):记录体(body)的字节数,不含 length 字段自身;记录总大小 = varint(length) 的字节数 + length。
  • attributes (int8):当前 8 位全部未使用,源码 DefaultRecord.java 写注释 "there are no used record attributes at the moment"。
  • timestampDelta (varlong):与批的 baseTimestamp 的差值,实际时间戳 =baseTimestamp + timestampDelta;若为 LogAppendTime 时间戳类型,读取时会直接用 broker 的 log append time 覆盖。
  • offsetDelta (varint):与批的 baseOffset 的差值,实际偏移 =baseOffset + offsetDelta;同时幂等/事务场景下生产者 sequence 也由baseSequence + offsetDelta推导(见DefaultRecord.readFrom中的DefaultRecordBatch.incrementSequence)。
  • keyLength / key:key 长度与字节内容;key 为 null 时 keyLength 编码为 -1(varint 的 -1),而非 0——这是"无 key"与"空 key"的区分。
  • valueLength / value:同上,null 时编码为 -1;Kafka 的 tombstone(墓碑消息)正是 value 为 null 的记录,用于 log compaction 删除逻辑。
  • headersCount / Headers:头部数量(varint)及头部数组。

正是由于这种逐字节紧凑编码,DefaultRecord中定义了固定开销常量MAX_RECORD_OVERHEAD = 21(注释:5 字节 length + 10 字节 timestamp + 5 字节 offset + 1 字节 attributes,不含 key/value/headers),用于估算一条记录的最大额外开销,供内存分配与批量大小估算使用。

记录头部(Record Header)

单个 Headers 条目(即一条消息 Header)的磁盘格式:

headerKeyLength: varint headerKey: String headerValueLength: varint Value: byte[]

文档明确了两个关键约束:

  1. header 的 key 保证非 null;而header 的 value 可以为 null
  2. Headers 的顺序在生产与消费过程中被完整保留——这意味着你可以在 Header 中存放有序的元数据序列(如链路追踪的 trace ID 列表),Kafka 不会打乱它们。

编码细节在 DefaultRecord.java 的writeTo中可见:header key 按 UTF-8 编码后以 varint 长度前缀写入,header value 以字节数组写入(null 时 varint 写 -1);读取时若 key 长度或 headers 数量为负会抛出InvalidRecordException,而 header 数量超过缓冲区剩余字节同样视为非法结构。

客户端的用户侧实现为 clients/src/main/java/org/apache/kafka/common/header/internals/RecordHeader.java:构造函数RecordHeader(String key, byte[] value)显式调用Objects.requireNonNull(key, "Null header keys are not permitted")拒绝 null key;key()value()采用双重检查锁 + 字段置空的懒加载策略,将ByteBuffer形态惰性转换为 String/byte[],避免不必要的拷贝。Record接口规定 magic 1 及以下版本headers()恒返回空数组——Headers 是 magic 2 才有的能力。

varint / varlong 编码

文档明确:Kafka 使用与Protobuf 相同的 varint 编码(每个字节 7 位有效负载 + 1 位续位标记,小整数占用更少字节),varlong 即变长 64 位整数。DefaultRecord通过 ByteUtils 的writeVarint/writeVarlong/readVarint/readVarlong读写,headers 数量同样以 varint 编码。

源码对照:从接口到实现的完整链路

梳理clients模块中与消息结构直接相关的核心类型:

  • RecordBatch.java:批抽象,定义 magic 常量(V0/V1/V2)、NO_PRODUCER_ID = -1NO_PRODUCER_EPOCH = -1NO_SEQUENCE = -1NO_PARTITION_LEADER_EPOCH = -1等"无值"哨兵,以及streamingIterator(BufferSupplier)——它延迟解压记录流,直到调用方真正请求下一条记录时才解压,且文档提示对 LZ4 这类需要 64KB 缓冲的解压算法,复用缓冲的 supplier 对迭代性能影响显著。
  • DefaultRecordBatch.java:magic 2 批的具体读写实现。
  • DefaultRecord.java:magic 2 单条记录实现,内含writeTo(序列化)、readFrom(反序列化并校验)、readPartiallyFrom(跳过 key/value/headers 的轻量读取,用于不需要完整负载的场景)与sizeInBytes系列估算方法。
  • Record.java:单条记录抽象(offset/sequence/timestamp/key/value/headers)。
  • MemoryRecords.java:内存中的批集合容器,提供batches()迭代器与多组withRecords/withIdempotentRecords/withTransactionalRecords工厂。
  • LegacyRecord.java:magic 0/1 旧格式实现。

旧消息格式(Old Message Format)

在 Kafka 0.11 之前,消息以**消息集(message sets)**的形式传输和存储,对应 magic 值 0 与 1(源码中的MAGIC_VALUE_V0/MAGIC_VALUE_V1)。旧格式没有独立的批头概念:每条消息各自携带完整的元数据,且RecordBatch注释指出——旧格式下若未启用压缩,一个"批"通常只含一条记录;压缩时一个批才可能包含多条记录。而 magic 2+ 的新格式无论是否压缩,一个批普遍包含多条记录,并且支持 Headers、事务与幂等语义。Record接口的若干方法(hasMagicisCompressedhasTimestampTypeheaders)都保留了旧格式的分支语义(如 magic < 2 时headers()返回空数组),以保证对旧数据段的兼容读取。新版客户端消费旧格式日志时,Kafka 会在读取路径上完成格式升级/转换。

总结

  • 一条 Kafka 消息 = 变长头部 + 不透明 key + 不透明 value;key/value 的序列化职责完全交由应用层,Kafka 内核只按字节处理。
  • 消息永远以记录批为单位写入磁盘,批头承载压缩算法、时间戳类型、事务/控制批标记、producer 状态等批量级元数据,并用 CRC-32C 校验 attributes 之后的全部字节。
  • 批内单条记录用 varint/varlong + delta 的紧凑编码描述 timestamp、offset 与 sequence,Headers 以"key 非空、value 可空、顺序保留"的约束承载应用元数据。
  • 控制批(ABORT/COMMIT)是事务协议与 KRaft 的底层信令,对应用完全透明。
  • 日志压缩为恢复生产者状态会保留批的首尾 offset/sequence,可能产生空批并改变 baseTimestamp,理解这一行为是排查幂等/事务生产异常的前提。

若需深入格式细节,可直接阅读 docs/implementation/message-format.md 原文,以及clients/src/main/java/org/apache/kafka/common/record/internal/目录下的完整实现与配套测试。

【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询