
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.jsonMetadataRequest.json / MetadataResponse.jsonProduceRequest.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 引入 CurrentLeaderEpochKIP-320、v12 引入 flexible versions、v13 用 topic ID 取代 topic 名称KIP-516、v15 引入 ReplicaStateKIP-903、v17 支持目录 IDKIP-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 起在没有 KIPKafka 改进提案的情况下不再允许删除对旧消息版本的支持因此不要随意抬高任何消息的版本支持区间下界。此外规格文件还可以用deprecatedVersions标注已废弃但仍兼容的版本例如 FetchRequest.json 中deprecatedVersions: 0-3、MetadataRequest.json 中deprecatedVersions: 0-3。请求/响应文件中的顶层元数据字段还包括字段含义示例值来自 FetchRequest.jsonapiKeyAPI 唯一标识1type消息类型request/responserequestlisteners该消息适用的监听器类型[zkBroker, broker, controller]name生成类的名称FetchRequestvalidVersions支持的版本范围0-17deprecatedVersions已废弃但兼容的版本范围0-3flexibleVersions启用灵活版本序列化的版本范围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 TypesKafka 消息协议提供以下原始字段类型类型说明bool布尔值true 或 falseint88 位有符号整数int1616 位有符号整数uint1616 位无符号整数int3232 位有符号整数uint3232 位无符号整数int6464 位有符号整数float64双精度浮点数IEEE 754stringUTF-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: 12ReplicaState结构使用tag: 1, taggedVersions: 15ReplicaDirectoryId使用tag: 0, taggedVersions: 17。4.4 灵活版本Flexible VersionsKafka 的序列化机制随版本演进不断改进包含这些改进的消息版本被称为灵活版本flexible versions。在灵活版本中string、array、bytes等变长字段以更节省空间的方式序列化。这些新的序列化类型以compact开头例如COMPACT_STRING是STRING的高效形式COMPACT_ARRAYOF、COMPACT_BYTES同理。规格文件通过顶层flexibleVersions属性声明哪些版本启用了灵活序列化例如FetchRequest.jsonflexibleVersions: 12MetadataRequest.jsonflexibleVersions: 9标记字段只能添加到灵活版本中tagged fields can only be added to flexible message versions这是两者之间的重要耦合关系。五、序列化与反序列化read / write / size5.1 序列化Message#writeMessage#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浮点数float640布尔值boolfalse字符串string空字符串字节bytes空字节数组UUID零 UUIDzero uuidrecordsnull数组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、LastFetchedEpochignorable: false、ForgottenTopicsDataignorable: 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-iJSON 规格文件的输入目录--typeclass-generators-t类型类生成器可多个--message-class-generators-m消息类生成器可多个处理过程见processDirectories方法用 JacksonObjectMapper解析输入目录下所有*.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.gradlegroup-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唯一、typerequest/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。灵活版本与标记字段配套标记字段tagtaggedVersions开放式区间只能加在 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),仅供参考