资讯动态

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

发布时间:2026/9/10 13:28:36 来源:尧图企业网站定制
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它继承IterableRecord同时暴露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 值为 2baseOffset: 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 1212 8 字节 baseOffset 4 字节 batchLength 自身。partitionLeaderEpoch (int32)分区 leader 的纪元号。magic (int8)当前取值为 2对应源码中的CURRENT_MAGIC_VALUE MAGIC_VALUE_V2另存在MAGIC_VALUE_V0 0、MAGIC_VALUE_V1 1。crc (uint32)覆盖从 attributes 到批末尾的所有字节即 CRC 之后的所有内容使用 CRC-32CCastagnoli多项式。CRC 位于 magic 之后因此客户端必须先解析 magic 字节才能决定如何解释 batchLength 与 magic 之间的字节。partitionLeaderEpoch 字段不参与 CRC 计算——这是刻意为之该字段由 broker 在接收每个批时赋值若不排除它每次赋值都需重算 CRC。attributes 位域压缩、时间戳类型、事务与控制批attributes 的 16 位被细分为多个语义位位含义取值说明bit 0~2压缩算法0: 无压缩1: gzip2: snappy3: lz44: zstdbit 3timestampType区分 CreateTime 与 LogAppendTimebit 4isTransactional1 表示该批属于某个事务bit 5isControlBatch1 表示控制批bit 6hasDeleteHorizonMs1 表示 baseTimestamp 被设置为压缩删除边界delete horizonbit 7~15未使用保留位压缩算法常量与编码对应关系在 clients/src/main/java/org/apache/kafka/common/record/internal/CompressionType.java 中定义。当启用压缩时压缩后的记录数据被直接序列化在 recordsCount 字段之后——也就是说压缩针对的是整批记录的字节流而不是逐条记录独立压缩。与压缩相关的生产者侧实践生产者端的compression.type配置可取none、gzip、snappy、lz4、zstd最终会写入上述 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 6hasDeleteHorizonMs从而告诉压缩器何时可以安全地清除这些记录。对应的 API 是RecordBatch#deleteHorizonMs()RecordBatch.java返回OptionalLong若批的首时间戳不是 delete horizon 则返回空。控制批Control Batch事务与 KRaft 的底层信令控制批Control Batch是 magic 2 引入的特殊批类型其中只包含一条称为控制记录control record的记录。控制记录绝不会返回给应用程序其用途有二消费者侧过滤消费者借助控制记录过滤掉已中止abort的事务消息KRaft 协议元数据控制批被 KRaft 共识协议用于承载内部协议元数据。控制记录的 key 遵循如下 schemaversion: int16 (current version is 0) type: int16 (the control record types are in the table below)常规 topicregular topics当前定义的控制记录类型如下TypeNameDescription0ABORT标记一个事务已中止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 / keykey 长度与字节内容key 为 null 时 keyLength 编码为 -1varint 的 -1而非 0——这是无 key与空 key的区分。valueLength / value同上null 时编码为 -1Kafka 的 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[]文档明确了两个关键约束header 的 key 保证非 null而header 的 value 可以为 nullHeaders 的顺序在生产与消费过程中被完整保留——这意味着你可以在 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 keykey()与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 -1、NO_PRODUCER_EPOCH -1、NO_SEQUENCE -1、NO_PARTITION_LEADER_EPOCH -1等无值哨兵以及streamingIterator(BufferSupplier)——它延迟解压记录流直到调用方真正请求下一条记录时才解压且文档提示对 LZ4 这类需要 64KB 缓冲的解压算法复用缓冲的 supplier 对迭代性能影响显著。DefaultRecordBatch.javamagic 2 批的具体读写实现。DefaultRecord.javamagic 2 单条记录实现内含writeTo序列化、readFrom反序列化并校验、readPartiallyFrom跳过 key/value/headers 的轻量读取用于不需要完整负载的场景与sizeInBytes系列估算方法。Record.java单条记录抽象offset/sequence/timestamp/key/value/headers。MemoryRecords.java内存中的批集合容器提供batches()迭代器与多组withRecords/withIdempotentRecords/withTransactionalRecords工厂。LegacyRecord.javamagic 0/1 旧格式实现。旧消息格式Old Message Format在 Kafka 0.11 之前消息以**消息集message sets**的形式传输和存储对应 magic 值 0 与 1源码中的MAGIC_VALUE_V0/MAGIC_VALUE_V1。旧格式没有独立的批头概念每条消息各自携带完整的元数据且RecordBatch注释指出——旧格式下若未启用压缩一个批通常只含一条记录压缩时一个批才可能包含多条记录。而 magic 2 的新格式无论是否压缩一个批普遍包含多条记录并且支持 Headers、事务与幂等语义。Record接口的若干方法hasMagic、isCompressed、hasTimestampType、headers都保留了旧格式的分支语义如 magic 2 时headers()返回空数组以保证对旧数据段的兼容读取。新版客户端消费旧格式日志时Kafka 会在读取路径上完成格式升级/转换。总结一条 Kafka 消息 变长头部 不透明 key 不透明 valuekey/value 的序列化职责完全交由应用层Kafka 内核只按字节处理。消息永远以记录批为单位写入磁盘批头承载压缩算法、时间戳类型、事务/控制批标记、producer 状态等批量级元数据并用 CRC-32C 校验 attributes 之后的全部字节。批内单条记录用 varint/varlong delta 的紧凑编码描述 timestamp、offset 与 sequenceHeaders 以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),仅供参考

读完文章,也想定制专属网站?

尧图设计师 24 小时内与您沟通定制方案

免费获取报价