
Apache Kafka 通信协议详解请求/响应格式、API 版本协商与客户端实现指南【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka本文是 Kafka 仓库中 docs/design/protocol.md 协议文档的技术导读系统讲解 Kafka 的 TCP 线协议wire protocol从网络模型、分区引导bootstrapping、分区策略、批处理与兼容性策略到 API 版本协商、SASL 认证流程、请求/响应的二进制格式BNF 文法、错误码与 ApiKey 约定并针对第三方客户端给出成员 ID 格式的官方建议。读完本文你将具备从零实现一个 Kafka 客户端或深入理解现有客户端源码所需的完整协议知识。概览一个请求-响应风格的二进制 TCP 协议Kafka 的协议本质上是基于 TCP 的二进制协议所有 API 都被定义为请求-响应消息对request-response message pairs所有消息都是大小定界size delimited的由若干基本类型primitive types拼接而成。协议层面不要求连接/断开的握手客户端建立 socket 连接后即可写入一串请求消息并依次读回对应的响应消息。关于连接的几个关键事实持久连接收益更高TCP 握手成本虽不高但维持长连接、在一条连接上复用多次请求能摊薄握手开销因此生产客户端通常采用连接池复用的方式。需要连接多个 broker数据被分区存放客户端必须与持有目标分区数据的多个 broker 分别建连但同一客户端实例通常不需要对单个 broker 维持多条连接即无需连接池化。单连接上的严格顺序保证服务端保证在一条 TCP 连接上请求按发送顺序被处理、响应也按该顺序返回。为此broker 的请求处理在单连接上只允许一个 in-flight 请求。客户端可且应该使用非阻塞 IO 实现请求流水线pipelining即使前一个请求的响应未返回也可以继续发送后续请求未完成的请求会缓冲在操作系统 socket 缓冲区中。请求大小上限服务端对请求大小有可配置的最大限制socket.request.max.bytes默认 104857600 字节超过该上限的请求会导致 socket 被直接断开。该协议的所有请求都由客户端发起并除非特别注明都会得到对应的响应消息。分区与引导Partitioning and bootstrappingKafka 是分区系统并非所有服务器都持有完整数据集。主题topic被拆分为预定义数量的分区 P每个分区按复制因子 N 复制每个分区本身就是一个编号为 0, 1, ..., P-1 的有序提交日志commit log。谁来决定数据落到哪个分区分区分配完全由客户端控制。broker 不强制任何消息该发往哪个分区的语义发布消息时客户端直接把消息寻址到某个具体分区拉取消息时也从某个具体分区拉取。如果两个客户端希望使用相同的分区方案就必须使用相同的key 到 partition映射算法。一个重要的约束是发布/拉取请求必须发送给当前作为该分区 leader 的 broker。这一条件由 broker 强制执行——向错误 broker 发送针对某分区的请求会返回NotLeaderForPartition错误码在 Errors.java 中对应NOT_LEADER_OR_FOLLOWER(6)其语义为该请求仅面向 leader而当前 broker 不是该 topic-partition 的 leader。元数据引导流程客户端如何知道集群里有哪些主题、它们有哪些分区、这些分区当前由哪些 broker 托管这些信息是动态的无法靠静态配置文件解决。为此所有 Kafka broker 都能应答元数据请求MetadataRequest描述集群当前状态有哪些主题、主题有哪些分区、各分区 leader 是谁、各 broker 的主机与端口。因此客户端的引导思路是先找到任意一个 broker由它告知客户端集群中其余 broker 及其托管的分区。由于这第一个 broker 本身可能宕机官方建议客户端实现应配置2~3 个引导 URL可由用户使用负载均衡器或静态配置多个 Kafka 主机。客户端无需轮询感知集群变化它可以在实例化时拉取一次元数据并缓存直到收到元数据过期的错误。该错误有两种形式socket 错误——客户端无法与某个 broker 通信请求响应中的错误码——该 broker 已不再托管客户端所请求的分区。对应的标准流程是循环尝试引导 Kafka URL 列表直到找到可连接的 broker拉取集群元数据处理 fetch / produce 请求根据目标 topic/partition 将请求定向到正确的 broker收到相应错误后刷新元数据并重试。从实现角度看元数据请求的版本演进在 MetadataRequest.jsonapiKey3中有完整注释v0 时空数组表示请求所有主题的元数据v1 时空数组改为不请求任何主题、null 数组表示请求所有主题v9 起成为第一个 flexible 版本v11 弃用IncludeClusterAuthorizedOperations字段改由 DescribeCluster API 提供即 KIP-700。这些注释正是协议文档动态元数据 缓存失效刷新机制在消息模式定义中的落地体现。分区策略Partitioning Strategies正如上文所述消息到分区的分配由生产客户端控制。分区在 Kafka 中承担两个目的负载均衡在 broker 之间均衡数据与请求负载语义分区在消费者进程之间划分处理任务同时允许在分区内保持本地状态与顺序。实际场景中你可能只关心其一或两者兼顾。实现层面的常见策略简单轮询round robin客户端将请求在所有 broker 上轮询分发适合纯粹为了负载均衡的场景。随机单分区在生产者数量远多于 broker的环境中让每个客户端随机选择一个分区发布。该策略产生的 TCP 连接数要少得多。语义分区按 key 哈希使用消息携带的某个 key 来分配分区。例如处理点击流时按用户 ID 分区使同一用户的所有数据都流向同一个消费者。实现方式是对 key 做哈希再用哈希值选择目标分区。批处理BatchingKafka 的 API 刻意鼓励把小数据批量聚合成大块来提升效率这被官方明确视为非常显著的性能收益发送消息的 API 与拉取消息的 API始终以消息序列sequence of messages为单位工作而非单条消息一个聪明的客户端可以实现异步模式把逐条发送的消息攒成更大的批次再发出Kafka 更进一步允许跨主题、跨分区的批量一次 produce 请求可包含追加到多个分区的数据一次 fetch 请求可一次性从多个分区拉取数据。当然客户端实现者也可以选择忽略批处理逐条发送——协议并不强制。兼容性API 版本协商Kafka 采用**双向客户端兼容策略**bidirectional client compatibility新客户端可以连接旧服务端旧客户端也可以连接新服务端。这样用户可以在不中断服务的情况下单独升级客户端或服务端。由于协议随时间演变客户端与服务端必须就线上消息的 schema达成一致手段就是API 版本化API versioning每次请求发送前客户端会先发送API key与API version两个 16 位数字两者组合起来唯一标识后续消息的 schema客户端应支持一个版本区间与特定 broker 通信时使用双方都支持的最高版本并在请求中标注服务端会拒绝它不支持的版本并且总是按照请求中携带的版本所对应的格式返回响应升级路径的设计意图是新特性先部署到服务端旧客户端暂时不使用待新客户端逐步上线后再逐渐利用新特性。唯一的例外是检索支持的 API 版本时服务端可能以不同的版本响应。KIP-482标记字段Tagged Fields协议的 flexible 版本支持KIP-482 引入的 tagged fields在不递增版本号的情况下向请求添加字段为消息 schema 的演进提供了额外途径。关键特性未设置的 tagged field不占任何空间因此对极少使用的字段做成 tagged field 比放进强制 schema 更省空间但 tagged fields 会被不认识它们的接收方静默忽略——如果这并非发送方想要的行为例如需要对方强制校验则可能带来隐患此时递增版本号可能更合适。检索支持的 API 版本ApiVersions自 0.10.0.0 起见 KIP-35broker 会暴露自身支持的各 API 版本信息客户端据此选择双方都支持的最高版本若不存在交集则应向用户报告错误。客户端获取版本的推荐流程与 broker 建立连接后若启用 SSL则在 SSL 握手完成后发送ApiVersionsRequestbroker 收到后无论当前认证状态如何都会返回其支持的 ApiKeys 与版本全列表。几点注意若担心此行为泄露 broker 版本信息可改用带客户端认证的 SSL因为 SSL 客户端认证发生在ApiVersionRequest之前版本早于 0.10.0.0 的 broker 不支持该 API会忽略请求或直接关闭连接若客户端请求的ApiVersionsRequest版本不被支持客户端版本超前且 broker 版本 ≥ 2.4.0broker 会返回 version 0 的ApiVersionsResponse错误码设为UNSUPPORTED_VERSION见 Errors.java 中UNSUPPORTED_VERSION(35)并填充 broker 支持的ApiVersionsRequest版本客户端据此重试用双方支持的最高版本再次发起请求见 KIP-511在 broker 上收集与暴露客户端名称和版本若某个 API 存在多个双方都支持的版本客户端应使用 broker 与自己都支持的最新版本协议版本的弃用通过将某 API 版本在协议文档中标记为 deprecated 完成重要从 broker 获取到的支持的 API 版本仅对获取该信息的这条连接有效。连接断开后客户端应重新获取因为期间 broker 可能已被升级/降级。ApiVersionsRequest的模式定义见 ApiVersionsRequest.jsonapiKey18当前有效版本为0-5v3 起为第一个 flexible 版本并新增ClientSoftwareName/ClientSoftwareVersionKIP-511 的实现载体v4 修复了 KAFKA-17011v5 引入ClusterId与NodeId校验以及REBOOTSTRAP_REQUIRED错误KIP-1242。SASL 认证序列SASL 认证遵循以下流程可选客户端先发送ApiVersionsRequest获取 broker 支持的请求版本范围客户端发送携带 SASL 机制的SaslHandshakeRequest。若请求的机制未在服务端启用服务端会返回支持的机制列表并关闭连接若已启用则返回成功响应并继续 SASL 认证执行真正的 SASL 认证若SaslHandshakeRequest为v0一系列 SASL 客户端/服务端 token 作为不透明包直接发送不包裹 Kafka 协议头若为v1改用SaslAuthenticate请求/响应实际 SASL token 被包裹在 Kafka 协议中broker 最终消息中的错误码将指示认证成功或失败认证成功后后续数据包作为 Kafka API 请求处理否则关闭客户端连接。为保证与 0.9.0.x 客户端的互操作若服务端收到的第一个数据包不是合法 Kafka 请求则将其当作 SASL/GSSAPI 客户端 token 处理从该包开始执行 SASL/GSSAPI 认证跳过上述前两步。协议原始类型Protocol Primitive Types协议由以下基本类型构建仓库generator模块会基于消息 JSON 定义生成对应文档见clients/src/main/resources/common/message/下各*.json与 generator 源码固定宽度整数int8、int16、int32、int64变长整数varint、varlong无符号变长整数unsigned varint用于紧凑编码字符串与字节数组string/nullable string/bytes/nullable bytes均以长度前缀int16开头数组ARRAY以元素个数int32开头记录集合records/nullable records存储消息批次record batches布尔booleanint80 或 1UUID16 字节定长。Flexible 版本compact encoding标记为 flexible 的消息版本其变长字段改用紧凑编码——数组用COMPACT_ARRAY而非ARRAY字符串与字节数组用对应的COMPACT_*类型。紧凑编码将长度存储为无符号变长整数unsigned varint而非定宽整数并且每个请求/响应组件末尾都带一个tagged-fields 区段。这使得字段可以零开销地按需出现正是 KIP-482 标记字段的编码基础。请求/响应格式文法BNF 阅读指南下面给出的是请求与响应二进制格式的精确上下文无关文法BNF。BNF 刻意写得不够紧凑目的是让名称可读。文法规则一串产生式production表示拼接concatenation多个候选产生式之间用|分隔可用括号分组顶层定义总是最先给出子部分缩进表示。通用请求与响应结构所有请求与响应都源自以下文法RequestOrResponse Size (RequestMessage | ResponseMessage) Size int32字段描述message_size后续请求/响应消息的字节大小。客户端先读取这个 4 字节的整数 N再读取并解析随后的 N 字节请求内容请求与响应头不同的请求/响应版本对应不同版本的头headers头版本与各 API 消息描述一并给出。整体消息结构为Message RequestOrResponseHeader BodyRequestOrResponseHeader是版本化的请求或响应头Body是消息特有的主体。所有 API 消息、错误码与 ApiKey 的完整列表由generator根据clients/src/main/resources/common/message/下的 JSON 定义生成协议文档中通过 include-html 嵌入对应生成产物。Record Batch消息批次record batch格式的完整描述见 docs/implementation/messages.md消息由变长头部、变长不透明 key 字节数组与变长不透明 value 字节数组组成。key 与 value 保持不透明是刻意设计——序列化库的选择权留给具体应用RecordBatch接口本质上是消息的迭代器并提供面向 NIOChannel的批量读写专用方法。常量错误码与 Api Keys错误码Error Codes协议使用数字错误码指示服务端发生了什么问题客户端可将其翻译为异常或对应的错误处理机制。错误码全表由 generator 生成源码定义集中在 Errors.java例如NOT_LEADER_OR_FOLLOWER(6)——请求仅面向 leader/follower而当前 broker 不是该 topic-partition 的 leader协议文档中写作NotLeaderForPartitionUNSUPPORTED_VERSION(35)——请求的 API 版本不受支持ApiVersions 版本回退机制依赖该错误码。Api Keys每个请求类型在ApiKey字段中对应一个数字编码。源码中的完整枚举见 ApiKeys.java它直接映射到ApiMessageType中各消息的定义。列举仓库当前版本中的代表性 ApiKey完整列表以 ApiMessageType 生成的枚举为准请求类型说明PRODUCE生产消息批量追加到分区FETCH拉取消息可跨多分区批量LIST_OFFSETS列出分区可用偏移量METADATA获取集群元数据引导核心 APIOFFSET_COMMIT / OFFSET_FETCH提交/读取消费偏移量FIND_COORDINATOR定位 group 或事务协调者JOIN_GROUP / HEARTBEAT / SYNC_GROUP / LEAVE_GROUP经典消费组协议CONSUMER_GROUP_HEARTBEAT / CONSUMER_GROUP_DESCRIBE新版消费组协议SASL_HANDSHAKE / SASL_AUTHENTICATESASL 认证API_VERSIONS检索支持的 API 版本CREATE_TOPICS / DELETE_TOPICS / CREATE_PARTITIONS管理类操作INIT_PRODUCER_ID / ADD_PARTITIONS_TO_TXN / END_TXN事务DESCRIBE_ACLS / CREATE_ACLS / DELETE_ACLS访问控制VOTE / BEGIN_QUORUM_EPOCH / END_QUORUM_EPOCHKRaft 控制器选举raft 协议ALLOCATE_PRODUCER_IDS生产者 ID 分配BROKER_REGISTRATION / BROKER_HEARTBEATKRaft 模式下 broker 注册与心跳GET_TELEMETRY_SUBSCRIPTIONS / PUSH_TELEMETRY客户端遥测为什么不用 HTTP / XMPP / Protobuf协议设计的哲学协议文档还回答了几个灵魂拷问这些决策直接塑造了 Kafka 的实现为什么不用 HTTP最好的理由是客户端实现者可以利用 TCP 的高级特性——请求多路复用multiplex、同时轮询多条连接等此外官方坦言许多语言的 HTTP 库质量参差不齐。为什么不支持多种协议历史经验表明新特性若要在多个协议实现间移植将极难添加和测试。多数用户并不把多协议当作特性他们只想要自己语言里一个可靠好用的客户端。为什么不用 XMPP/STOMP/AMQP 等现成协议协议在很大程度上决定了实现形态若无法掌控协议就无法实现当前的分布式消息能力。Kafka 团队的信念是在提供真正分布式消息系统这件事上可以做得比现有系统更好而这需要构建一个工作方式不同的东西。为什么不用 Protocol Buffers / Thrift这些工具擅长管理海量序列化消息但 Kafka 只有少数几种消息跨语言支持参差不齐更重要的是二进制日志格式与线协议之间的映射需要精细管理这是这些系统做不到的。此外Kafka 偏好显式的 API 版本化与校验而非把新值推断为 null的隐式演进因为前者能对兼容性做更精细的控制。第三方客户端建议Member ID 格式Kafka 客户端参与组协议例如ConsumerGroupHeartbeatRPC时需要生成一个member ID向 broker 标识自己。虽然协议并不强制其格式官方强烈建议使用base64 编码的 UUID作为 member ID采用URL-safe base64编码不含或/字符去掉连字符——最终字符串应为连续的字母数字序列例如abc123def456。示例标准 UUID如00000000-0000-0000-0000-000000000000应转换为类似YzYxNjQ4OTItZDE1Mi00Y2E4LWIyNzUtYmIwMzAwMDAwMDAw的 URL-safe base64 字符串说明性示例实际编码取决于 UUID 字节。需要强调的是这是一个强建议而非强约束协议不会拒绝偏离该格式的 member ID。进一步阅读docs/design/protocol.md协议指南原文本文所依据的源文档docs/implementation/messages.md消息与 RecordBatch 格式详解clients/src/main/resources/common/message/所有请求/响应消息的 JSON 模式定义协议的真实数据源clients/src/main/java/org/apache/kafka/common/protocol/ApiKeys.java 与 Errors.javaApiKey 与错误码枚举generator由消息 JSON 生成协议文档与 Java 代码的代码生成器。【免费下载链接】KafkaApache Kafka - A distributed event streaming platform项目地址: https://gitcode.com/GitHub_Trending/kafka4/kafka创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考