news 2026/9/19 8:55:54

Apache Kafka 消息协议定义体系详解:基于 JSON 规格文件的消息生成机制

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Kafka 消息协议定义体系详解:基于 JSON 规格文件的消息生成机制

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 协议由"请求-响应"构成:客户端向服务器发送请求以获取响应。协议规定了两条基本规则:

  1. 每个请求由 16 位整数"apiKey"唯一标识。例如在 FetchRequest.json 中"apiKey": 1,在 MetadataRequest.json 中"apiKey": 3。响应的 apiKey 永远与请求一致。
  2. 每个消息有独立的 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)
apiKeyAPI 唯一标识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),包含readwritesize等方法,这部分内容在"序列化与反序列化"一节详述。

四、字段(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)、tagtaggedVersions(标记字段)、ignorable(可忽略字段)、entityType(语义类型标注,如brokerIdtopicName)、嵌套的fields(子结构)。其中"versions": "0-14""versions": "15+"的配合(ReplicaId 被 ReplicaState 取代)正是"移除旧字段、新增新字段"的典型写法。

4.1 字段类型(Field Types)

Kafka 消息协议提供以下原始字段类型:

类型说明
bool布尔值,true 或 false
int88 位有符号整数
int1616 位有符号整数
uint1616 位无符号整数
int3232 位有符号整数
uint3232 位无符号整数
int6464 位有符号整数
float64双精度浮点数(IEEE 754)
stringUTF-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;而stringbytesuuidrecords、数组类型字段可以选择性地声明为可空。字段"可空"意味着序列化/反序列化代码准备处理该字段的 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 协议的扩展机制,允许向消息附加可选数据。标记字段可以出现在消息的根层级,也可以出现在消息内的任何结构(如嵌套结构体)中。

与必填字段不同,标记字段可以添加到已经存在的消息版本上,且旧版本服务器会忽略它们不理解的标记字段——这为协议演进提供了极大的灵活性。

使字段成为标记字段需要两步:

  1. 为字段设置tag(一个整数标识);
  2. 设置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)

在灵活版本中,stringarraybytes等变长字段以更节省空间的方式序列化。这些新的序列化类型以compact开头,例如COMPACT_STRINGSTRING的高效形式(COMPACT_ARRAYOFCOMPACT_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.ObjectSerializationCacheorg.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)
recordsnull
数组(array)空集合

字符串字段可以通过指定字面量字符串"null"把默认值设为 null(例如 FetchRequest.json 中ClusterId"default": "null")。注意:只有当字段的所有版本都可空时,才能把 null 指定为默认值

5.3 自定义默认值(Custom Default Values)

对于整数、布尔、浮点、字符串类型的字段,可以在 JSON 对象中添加default条目设置自定义默认值,它覆盖该类型的常规默认值。例如,可以让某个布尔字段默认值为true而非false

自定义默认值必须对字段类型有效:int16 字段的默认值必须是能装进 16 位的整数,以此类推。可以使用十六进制或八进制写法,只要分别以0x0开头即可(如 FetchRequest.json 中MaxBytes的默认值"0x7fffffff")。目前不能为 bytes 或数组字段设置自定义默认值。

自定义默认值的典型用途:当旧版本消息缺少某些信息时给出合理的兜底。例如,旧版本请求没有 timeout 字段,可以指定服务器假设这类请求的超时时间为 5000ms 或其他任意值,从而保持新旧版本行为一致。

5.4 可忽略字段(Ignorable Fields)

用旧或新格式写消息时,并非所有字段都会出现;接收方反序列化时会给缺失字段填上默认值。因此,如果源字段被设置为非默认值,这部分信息就会丢失

  • 某些情况下信息丢失可以接受(如 timeout 字段);
  • 某些情况下字段非常重要、不应丢弃(如改变请求整体含义的 "verify only" 布尔字段)。

默认行为是:信息丢失不被允许——如果被忽略的字段没有设置为默认值,消息序列化代码会抛出异常。如果某个字段的信息丢失是可以接受的,请为该字段设置"ignorable": true以关闭此检查;此时该字段可以在序列化时被静默省略。

真实例子:FetchRequest.json 中MaxBytesIsolationLevelSessionIdSessionEpochRackIdCurrentLeaderEpochLogStartOffsetReplicaDirectoryId等都标记了"ignorable": true(这些字段通常有明确的默认值语义,旧版本丢字段不影响正确性),而ClusterIdLastFetchedEpoch"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 列出的四类典型不兼容变更:

  1. 修改已发布的协议版本。已发布的协议版本必须被视为"既成事实";如果发现错误,应在新版本中修正,而不是改动既有版本。
  2. 重排已有字段。允许在已有字段之前或之后新增字段,但已有字段之间不应相互重排(因为字段顺序即线上字节顺序)。
  3. 改变已有字段的默认值。绝不能修改已存在字段的默认值,否则新旧客户端与服务器会对默认值产生分歧。
  4. 改变已有字段的类型。唯一的例外是:只要转换正确,原始类型数组可以改成包含相同数据的结构体数组。因为 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方法):

  1. 用 JacksonObjectMapper解析输入目录下所有*.json文件为MessageSpec对象;注意其配置开启了JsonParser.Feature.ALLOW_COMMENTS,这正是 README 所说"JSON 支持//注释"的实现位置(MessageGenerator.java);
  2. 依次调用各消息类生成器(如MessageDataGeneratorJsonConverterGenerator)生成对应*.java文件;
  3. 调用类型类生成器(如ApiMessageTypeGeneratorMetadataRecordTypeGenerator)汇总生成ApiMessageType.javaMetadataRecordType.java等全局注册文件;
  4. 清理输出目录中不再由任何规格文件生成的旧文件。

在构建层面,各模块通过 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 批准)。
  • 字段属性齐全:每个字段有nametypeversions;字符串/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),仅供参考

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/19 8:55:18

图吧工具箱深度解析:原理、下载、权限与WMI监控机制

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华
网站建设 2026/9/19 8:50:28

CMSIS-4老工程迁移指南:从结构尽调到AC6实战避坑

/* MD / 富文本中的 .toc(含博客园搬家等嵌套结构);.toc-box 在侧栏,不受影响 */#content_views .toc,/* 编辑器常在目录前后插入空 p(:empty 仍占 20px),一并去掉避免顶空隙 */#content_views.markdown_views > p:empty:has(+ .toc),#content_views.markdown_views …

作者头像 李华