Apache Kafka 消息协议定义体系详解:基于 JSON 规格文件的消息生成机制
【免费下载链接】kafkaMirror of Apache Kafka项目地址: https://gitcode.com/gh_mirrors/kafka31/kafka
本篇技术指南以 Kafka 仓库中 clients/src/main/resources/common/message/README.md 为骨架,系统讲解 Apache Kafka 如何用 JSON 文件定义客户端与服务器之间传输的消息协议(请求/响应结构、字段类型、版本演进与序列化规则),并结合仓库内的真实消息规格(如 FetchRequest、MetadataRequest)、生成器源码与 Gradle 构建任务,说明这套"声明式定义 + 代码自动生成"机制的工作方式。读完本文,你将掌握 Kafka 消息 JSON 文件的完整书写规范、各关键属性的语义(validVersions、flexibleVersions、nullableVersions、taggedVersions、mapKey、ignorable 等),以及如何从源码层面验证和理解协议变更的约束。
一、消息定义体系概述:从手写序列化到声明式生成
Kafka 的客户端(Producer、Consumer、AdminClient)与 Broker 之间通过一套二进制协议通信,这套协议定义了"请求(Request)"与"响应(Response)"的结构以及它们在网络上的序列化方式。
在 clients/src/main/resources/common/message/ 目录下,每一个 API 对应一个 JSON 规格文件,例如:
- FetchRequest.json / FetchResponse.json
- MetadataRequest.json / MetadataResponse.json
- ProduceRequest.json / ProduceResponse.json
这些 JSON 文件是 Kafka 消息协议的唯一事实来源。当 Kafka 被编译时,构建系统会把这些规格文件翻译成 Java 代码,用于读写消息;任何对 JSON 文件的修改都会触发对应生成代码的重新编译。这一点在 README 中有明确说明,并在构建脚本中得到印证(详见下文"生成流程"一节)。
这套体系取代了早期手写序列化代码的方案。从仓库中大量消息 JSON 文件的存在(共一百余个,涵盖 Produce、Fetch、JoinGroup、OffsetCommit、事务协调器、KRaft 共识、Share Group 等全部 API)可以看出,社区正持续把所有消息迁移到自动生成序列化/反序列化代码的轨道上。
需要特别指出的是:该 JSON 格式支持注释,注释以双斜杠//开头。这使得每个版本演进的原因都可以直接写在规格文件里,例如 FetchRequest.json 顶部就以注释形式记录了从 Version 0 到 Version 17 的完整演进历史(从 v2 引入 MaxBytes、v4 引入 IsolationLevel、v7 引入增量抓取会话、v9 引入 CurrentLeaderEpoch(KIP-320)、v12 引入 flexible versions、v13 用 topic ID 取代 topic 名称(KIP-516)、v15 引入 ReplicaState(KIP-903)、v17 支持目录 ID(KIP-853)等),极大提升了协议规格的可读性与可追溯性。
二、请求与响应:apiKey 与版本号的核心约定
Kafka 协议由"请求-响应"构成:客户端向服务器发送请求以获取响应。协议规定了两条基本规则:
- 每个请求由 16 位整数"apiKey"唯一标识。例如在 FetchRequest.json 中
"apiKey": 1,在 MetadataRequest.json 中"apiKey": 3。响应的 apiKey 永远与请求一致。 - 每个消息有独立的 16 位版本号。不同版本的消息其 schema(字段集合)可能不同;有时版本号递增但 schema 未变,这可能只是提示服务器以某种不同方式处理该消息。响应的版本号必须与对应请求的版本号一致。
每个请求或响应都有一个顶层字段validVersions,声明当前代码能够理解的协议版本范围。例如"validVersions": "0-2"表示支持版本 0、1、2。必须始终指明所支持的最高消息版本。
关于版本区间的下界,README 有一条重要提醒:目前唯一不再支持的旧版本是 MetadataRequest 与 MetadataResponse 的版本 0。自 KIP-97 起,在没有 KIP(Kafka 改进提案)的情况下不再允许删除对旧消息版本的支持,因此不要随意抬高任何消息的版本支持区间下界。此外,规格文件还可以用deprecatedVersions标注已废弃但仍兼容的版本,例如 FetchRequest.json 中"deprecatedVersions": "0-3"、MetadataRequest.json 中"deprecatedVersions": "0-3"。
请求/响应文件中的顶层元数据字段还包括:
| 字段 | 含义 | 示例值(来自 FetchRequest.json) |
|---|---|---|
apiKey | API 唯一标识 | 1 |
type | 消息类型:request/response | "request" |
listeners | 该消息适用的监听器类型 | ["zkBroker", "broker", "controller"] |
name | 生成类的名称 | "FetchRequest" |
validVersions | 支持的版本范围 | "0-17" |
deprecatedVersions | 已废弃但兼容的版本范围 | "0-3" |
flexibleVersions | 启用"灵活版本"序列化的版本范围 | "12+" |
三、MessageData 对象:一份数据,多个版本
基于 JSON 文件,Kafka 会为每个消息生成对应的MessageDataJava 对象,用于在 JVM 内存中保存请求和响应数据。
一个关键设计是:MessageData 对象本身不包含版本号,单个 MessageData 对象可以代表一个消息的所有版本。这使得业务代码可以用同一套代码路径处理所有版本的消息——发送时指定版本、读取时按版本解释字段,从而避免为每个版本编写一套数据类。
从 generator/src/main/java/org/apache/kafka/message/ 目录下的生成器源码可以看出,生成体系包含多种生成器组件,例如:
MessageDataGenerator:生成消息数据类;JsonConverterGenerator:生成 JSON 转换器;ApiMessageTypeGenerator:汇总所有消息的 apiKey 与版本信息;MetadataRecordTypeGenerator/MetadataJsonConvertersGenerator:面向元数据记录(KRaft 元数据日志)的生成器。
生成后的消息数据类实现了org.apache.kafka.common.protocol.Message/ApiMessage接口(见 MessageGenerator.java),包含read、write、size等方法,这部分内容在"序列化与反序列化"一节详述。
四、字段(Fields):顺序、版本与结构
每个消息包含一个字段数组fields,字段定义了消息中携带的数据。一般而言,字段具有**名称(name)、类型(type)和版本信息(versions)**三个属性。
字段顺序是协议契约的一部分:字段在消息定义中出现的顺序,就是它们在网络上被发送的顺序。调整已有字段之间的相对顺序属于不兼容变更(详见"不兼容变更"一节)。
在每个新消息版本中,可以增删字段:
- 新增字段:例如为某个消息创建新版本 3 时,可以用
"versions": "3+"声明该字段只在版本 3 及以后出现; - 移除字段:将字段的版本从
"0+"改为"0-2",表示它在版本 3 及以后不再出现。
以下是一个真实示例——FetchRequest.json 的字段定义节选:
{ "name": "ClusterId", "type": "string", "versions": "12+", "nullableVersions": "12+", "default": "null", "taggedVersions": "12+", "tag": 0, "ignorable": true, "about": "The clusterId if known. This is used to validate metadata fetches prior to broker registration." }, { "name": "ReplicaId", "type": "int32", "versions": "0-14", "default": "-1", "entityType": "brokerId", "about": "The broker ID of the follower, of -1 if this request is from a consumer." }, { "name": "MaxWaitMs", "type": "int32", "versions": "0+", "about": "The maximum time in milliseconds to wait for the response." }, { "name": "MaxBytes", "type": "int32", "versions": "3+", "default": "0x7fffffff", "ignorable": true, "about": "The maximum bytes to fetch. See KIP-74 for cases where this limit may not be honored." }, { "name": "Topics", "type": "[]FetchTopic", "versions": "0+", "about": "The topics to fetch.", "fields": [ ... ] }该示例集中展示了本节与后续小节涉及的大部分属性:versions(出现版本区间)、nullableVersions(可空版本区间)、default(自定义默认值,支持十六进制0x7fffffff)、tag与taggedVersions(标记字段)、ignorable(可忽略字段)、entityType(语义类型标注,如brokerId、topicName)、嵌套的fields(子结构)。其中"versions": "0-14"与"versions": "15+"的配合(ReplicaId 被 ReplicaState 取代)正是"移除旧字段、新增新字段"的典型写法。
4.1 字段类型(Field Types)
Kafka 消息协议提供以下原始字段类型:
| 类型 | 说明 |
|---|---|
bool | 布尔值,true 或 false |
int8 | 8 位有符号整数 |
int16 | 16 位有符号整数 |
uint16 | 16 位无符号整数 |
int32 | 32 位有符号整数 |
uint32 | 32 位无符号整数 |
int64 | 64 位有符号整数 |
float64 | 双精度浮点数(IEEE 754) |
string | UTF-8 字符串 |
uuid | 类型 4 不可变全局唯一标识符 |
bytes | 二进制数据 |
records | 记录集,例如内存中的 MemoryRecords |
除原始类型外还有数组类型(Array):以[]开头、以元素类型名结尾,例如[]Foo表示 "Foo 对象数组"。数组字段自带其元素对象的字段数组fields,用于描述包含对象的结构。真实例子如 FetchRequest.json 中的"type": "[]FetchTopic"与"type": "[]FetchPartition"(嵌套数组),以及 MetadataRequest.json 中的"type": "[]MetadataRequestTopic"。
关于各类型在网络上的具体序列化字节布局,可参见 Kafka 官方的协议文档(本仓库 README 中提到的 Kafka Protocol Guide)。
4.2 可空字段(Nullable Fields)
布尔、整数和浮点类型永远不可为 null;而string、bytes、uuid、records、数组类型字段可以选择性地声明为可空。字段"可空"意味着序列化/反序列化代码准备处理该字段的 null 值。
可空性通过nullableVersions属性声明。之所以把可空性实现为版本区间,是为了兼容 Kafka 中非常常见的模式:某个原本不可空的字段在后续版本中变为可空。最典型的例子是 MetadataRequest.json 中的Topics字段:
{ "name": "Topics", "type": "[]MetadataRequestTopic", "versions": "0+", "nullableVersions": "1+", ... }其语义(文件内注释也做了说明):版本 0 中空数组表示"请求所有 topic 的元数据";从版本 1 起空数组表示"不请求任何 topic 的元数据",而null 数组才表示"请求所有 topic 的元数据"——这就是可空版本区间从1+开始的直接原因。
使用约定:如果字段声明为不可空且出现在你正在使用的消息版本中,那么序列化前必须将其设置为非 null 值,否则会产生运行时错误。
4.3 标记字段(Tagged Fields)
标记字段(Tagged Fields)是 Kafka 协议的扩展机制,允许向消息附加可选数据。标记字段可以出现在消息的根层级,也可以出现在消息内的任何结构(如嵌套结构体)中。
与必填字段不同,标记字段可以添加到已经存在的消息版本上,且旧版本服务器会忽略它们不理解的标记字段——这为协议演进提供了极大的灵活性。
使字段成为标记字段需要两步:
- 为字段设置
tag(一个整数标识); - 设置
taggedVersions版本区间。
taggedVersions应当是开放式(open-ended)的——即只指定起始版本而不指定结束版本(如"12+")。你可以从某个具体消息版本中移除对某个标记字段的支持,但一旦某个 tag 被用于某种用途,就不能再复用于其他用途,否则会破坏兼容性。
真实示例:FetchRequest.json 中ClusterId字段使用"tag": 0, "taggedVersions": "12+",ReplicaState结构使用"tag": 1, "taggedVersions": "15+",ReplicaDirectoryId使用"tag": 0, "taggedVersions": "17+"。
4.4 灵活版本(Flexible Versions)
Kafka 的序列化机制随版本演进不断改进,包含这些改进的消息版本被称为灵活版本(flexible versions)。
在灵活版本中,string、array、bytes等变长字段以更节省空间的方式序列化。这些新的序列化类型以compact开头,例如COMPACT_STRING是STRING的高效形式(COMPACT_ARRAYOF、COMPACT_BYTES同理)。
规格文件通过顶层flexibleVersions属性声明哪些版本启用了灵活序列化,例如:
- FetchRequest.json:
"flexibleVersions": "12+" - MetadataRequest.json:
"flexibleVersions": "9+"
标记字段只能添加到灵活版本中(tagged fields can only be added to flexible message versions),这是两者之间的重要耦合关系。
五、序列化与反序列化:read / write / size
5.1 序列化(Message#write)
Message#write方法把消息写入缓冲区。实际写入哪些字段取决于调用write()时提供的版本号:当用较旧版本写入消息时,在该版本 schema 中尚不存在的字段会被省略。
因此,处理旧版本消息时,务必确认旧版 schema 包含了所有需要发送的数据。README 给出了明确的取舍原则:
- 可以接受省略的字段,例如 timeout 字段;
- 不能忽略会从根本上改变请求语义的字段,例如
validateOnly布尔值(它决定请求是校验还是真正执行)。
在序列化之前,常常需要知道消息会占用多少空间,此时可以调用Message#size方法。从 MessageGenerator.java 可以看到,生成代码中还会使用org.apache.kafka.common.protocol.ObjectSerializationCache与org.apache.kafka.common.protocol.MessageSizeAccumulator来辅助计算消息尺寸。
5.2 反序列化(Message#read)
消息对象通过Message#read方法反序列化,该方法会用新数据覆盖消息对象中的全部现有数据。
反序列化时,凡是在当前版本中不存在的字段都会被重置为默认值。各类型的默认值如下:
| 字段类型 | 默认值 |
|---|---|
| 整数(int8/int16/int32/int64 等) | 0 |
| 浮点数(float64) | 0 |
| 布尔值(bool) | false |
| 字符串(string) | 空字符串"" |
| 字节(bytes) | 空字节数组 |
| UUID | 零 UUID(zero uuid) |
| records | null |
| 数组(array) | 空集合 |
字符串字段可以通过指定字面量字符串"null"把默认值设为 null(例如 FetchRequest.json 中ClusterId的"default": "null")。注意:只有当字段的所有版本都可空时,才能把 null 指定为默认值。
5.3 自定义默认值(Custom Default Values)
对于整数、布尔、浮点、字符串类型的字段,可以在 JSON 对象中添加default条目设置自定义默认值,它覆盖该类型的常规默认值。例如,可以让某个布尔字段默认值为true而非false。
自定义默认值必须对字段类型有效:int16 字段的默认值必须是能装进 16 位的整数,以此类推。可以使用十六进制或八进制写法,只要分别以0x或0开头即可(如 FetchRequest.json 中MaxBytes的默认值"0x7fffffff")。目前不能为 bytes 或数组字段设置自定义默认值。
自定义默认值的典型用途:当旧版本消息缺少某些信息时给出合理的兜底。例如,旧版本请求没有 timeout 字段,可以指定服务器假设这类请求的超时时间为 5000ms 或其他任意值,从而保持新旧版本行为一致。
5.4 可忽略字段(Ignorable Fields)
用旧或新格式写消息时,并非所有字段都会出现;接收方反序列化时会给缺失字段填上默认值。因此,如果源字段被设置为非默认值,这部分信息就会丢失。
- 某些情况下信息丢失可以接受(如 timeout 字段);
- 某些情况下字段非常重要、不应丢弃(如改变请求整体含义的 "verify only" 布尔字段)。
默认行为是:信息丢失不被允许——如果被忽略的字段没有设置为默认值,消息序列化代码会抛出异常。如果某个字段的信息丢失是可以接受的,请为该字段设置"ignorable": true以关闭此检查;此时该字段可以在序列化时被静默省略。
真实例子:FetchRequest.json 中MaxBytes、IsolationLevel、SessionId、SessionEpoch、RackId、CurrentLeaderEpoch、LogStartOffset、ReplicaDirectoryId等都标记了"ignorable": true(这些字段通常有明确的默认值语义,旧版本丢字段不影响正确性),而ClusterId、LastFetchedEpoch("ignorable": false)、ForgottenTopicsData("ignorable": false)等则被显式声明为不可忽略,提示协议实现者这些字段的丢失是有害的。
六、Hash Sets:用 mapKey 提升查找效率
Kafka 中非常常见的模式是把消息数组中的元素载入 Map 或 Set 以便快速访问。消息协议通过mapKey概念支持这一点:
- 如果数组的某些元素字段被标注
"mapKey": true,整个数组将被当作**链式哈希集合(linked hash set)**而不是普通列表处理; - 集合中的元素可以用自动生成的
find函数以 O(1) 时间访问; - 集合元素的顺序仍然保持(插入序),新加入的条目总是排在最后。
从生成器源码 MessageGenerator.java 可以看到,mapKey 集合在生成代码中映射为org.apache.kafka.common.utils.ImplicitLinkedHashCollection(及其多重集合变体ImplicitLinkedHashMultiCollection),这正是 README 所说"linked hash set"的底层实现。
真实示例(ApiVersionsResponse.json):
{ "name": "ApiKey", "type": "int16", "versions": "0+", "mapKey": true, "about": "API keys ..." }七、不兼容变更(Incompatible Changes):必须避开的红线
避免对消息协议做不兼容变更是极其重要的。README 列出的四类典型不兼容变更:
- 修改已发布的协议版本。已发布的协议版本必须被视为"既成事实";如果发现错误,应在新版本中修正,而不是改动既有版本。
- 重排已有字段。允许在已有字段之前或之后新增字段,但已有字段之间不应相互重排(因为字段顺序即线上字节顺序)。
- 改变已有字段的默认值。绝不能修改已存在字段的默认值,否则新旧客户端与服务器会对默认值产生分歧。
- 改变已有字段的类型。唯一的例外是:只要转换正确,原始类型数组可以改成包含相同数据的结构体数组。因为 Kafka 协议不对结构做"装箱"(boxing),一个只含单个 int32 的结构体数组与 int32 数组在协议层是等价的。
这些约束与"版本演进"的总体哲学一致:协议变更只允许"向后兼容地增加",并通过 apiKey + 版本号 + 可空区间 + 标记字段 + flexible versions 这套组合拳来实现。仓库中每个规格文件的版本注释(如 FetchRequest.json 从 v0 到 v17 的逐版本说明)正是这种"增量为王、旧版冻结"实践的直接体现。
八、生成流程:JSON 如何变成 Java 代码
把 JSON 规格翻译为 Java 代码的核心入口是 MessageGenerator.java 的main方法。它接受以下命令行参数:
| 参数 | 缩写 | 含义 |
|---|---|---|
--package | -p | 生成代码所属的 Java 包名 |
--output | -o | 生成代码的输出目录 |
--input | -i | JSON 规格文件的输入目录 |
--typeclass-generators | -t | 类型类生成器(可多个) |
--message-class-generators | -m | 消息类生成器(可多个) |
处理过程(见processDirectories方法):
- 用 Jackson
ObjectMapper解析输入目录下所有*.json文件为MessageSpec对象;注意其配置开启了JsonParser.Feature.ALLOW_COMMENTS,这正是 README 所说"JSON 支持//注释"的实现位置(MessageGenerator.java); - 依次调用各消息类生成器(如
MessageDataGenerator、JsonConverterGenerator)生成对应*.java文件; - 调用类型类生成器(如
ApiMessageTypeGenerator、MetadataRecordTypeGenerator)汇总生成ApiMessageType.java、MetadataRecordType.java等全局注册文件; - 清理输出目录中不再由任何规格文件生成的旧文件。
在构建层面,各模块通过 Gradle 的processMessages任务驱动该生成器。以 build.gradle 中 metadata 模块的任务为例:
task processMessages(type:JavaExec) { mainClass = "org.apache.kafka.message.MessageGenerator" classpath = configurations.generator args = [ "-p", "org.apache.kafka.common.metadata", "-o", "src/generated/java/org/apache/kafka/common/metadata", "-i", "src/main/resources/common/metadata", "-m", "MessageDataGenerator", "JsonConverterGenerator", "-t", "MetadataRecordTypeGenerator", "MetadataJsonConvertersGenerator" ] inputs.dir("src/main/resources/common/metadata") ... } compileJava.dependsOn 'processMessages'类似的processMessages任务在 build.gradle(group-coordinator 模块,输入目录正是src/main/resources/common/message)以及 build.gradle、build.gradle、build.gradle、build.gradle、build.gradle 等处均有定义,分别服务于 clients(生成org.apache.kafka.common.message下的请求/响应类)、metadata、group-coordinator 等模块。可见整个 Kafka 代码库的协议实现都统一依赖这套"规格驱动生成"的流水线。
九、实战总结:一份合格的消息规格文件应满足什么
综合以上规则,当在 Kafka 仓库中阅读或编写一份消息 JSON 规格文件时,可以按以下清单自检:
- 顶层元数据完整:
apiKey唯一、type(request/response)正确、name与文件名一致、validVersions指明最高支持版本、flexibleVersions按需声明。 - 版本区间准确:新增字段用
N+、移除字段将区间上限改为N-1;不要抬高validVersions下界(自 KIP-97 起需 KIP 批准)。 - 字段属性齐全:每个字段有
name、type、versions;字符串/bytes/uuid/records/数组类型按需声明nullableVersions;语义关键字段标注entityType;信息丢失可接受时标注ignorable: true;需要用 O(1) 查找时给数组元素标注mapKey: true。 - 灵活版本与标记字段配套:标记字段(
tag+taggedVersions,开放式区间)只能加在 flexible 版本上,且 tag 一经使用不可复用。 - 默认值合法:自定义
default必须与字段类型匹配(支持0x/0前缀的十六/八进制),bytes 与数组字段不支持自定义默认;只有字段所有版本均可空时才能用"null"作为默认值。 - 不触犯兼容红线:不改已发布版本、不重排已有字段、不修改已有字段默认值与类型。
借助这套声明式协议定义体系,Kafka 得以在保证线上兼容的前提下持续演进:新增 API 只需新增 JSON 文件、扩展 API 只需递增版本并增量声明字段,而序列化/反序列化、版本分支、尺寸计算等繁琐且易错的代码全部由 generator 自动生成——这正是 Kafka 协议能支撑十余年大规模演进、横跨数十个版本仍保持强兼容性的底层机制之一。
【免费下载链接】kafkaMirror of Apache Kafka项目地址: https://gitcode.com/gh_mirrors/kafka31/kafka
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考