资讯动态

Apache Pulsar 二进制协议规范详解:从帧结构到 Lookup 的服务端通信原理

发布时间:2026/9/25 2:46:03 来源:尧图企业网站定制
消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载Apache Pulsar 的客户端与 Broker 之间并不依赖 HTTP 或 gRPC 等通用 RPC 框架而是采用了一套自研的二进制协议Binary Protocol。该协议专为消息中间件的高吞吐场景设计在支持确认acknowledgement、流控flow control、批量消息等完整功能的同时最大化传输与实现效率。本文基于 2.3.0 版官方文档与当前仓库的 Protobuf 定义、Java 客户端协议实现完整拆解这套协议的帧结构、消息元数据、各阶段交互命令Connect、Producer、Consumer、Lookup以及分区主题发现机制帮助读者理解 Pulsar 数据面的底层通信原理并具备自行抓包分析或实现类 Pulsar 客户端的能力。协议总览基于 Protobuf 的 BaseCommandPulsar 客户端和 Broker 之间交换的基本单元是命令command。所有命令都被封装为二进制的Protocol Buffersprotobuf消息其格式由仓库中的 PulsarApi.proto 文件完整定义该文件即文档所称的“Protobuf interface”。从源码结构看Java 侧生成的类位于org.apache.pulsar.common.api.proto包下客户端与 Broker 的帧序列化/反序列化则集中在 Commands.java 中完成。所有协议命令都包含在一个BaseCommandprotobuf 消息中。BaseCommand内部定义了一个Type枚举把每一种子命令Connect、Ping、Send、Subscribe……都声明为一个 optional 字段而一条BaseCommand消息只允许携带一个子命令——这正是它被设计成“oneof 式”结构的原因接收方先读取type字段即可 O(1) 地定位到对应子命令的解析分支。Commands.java 中的localCmd(BaseCommand.Type type)方法正是这一模式的典型实现先setType(type)再往对应子命令中填充字段。需要特别记住的一点不同 Producer 和 Consumer 的命令可以在同一条 TCP 连接上无限制地交错发送连接复用 / connection sharing。这就是为什么后续几乎所有命令都要携带producer_id/consumer_id和request_id来区分上下文与配对请求/响应。Framing协议帧结构protobuf 本身不提供任何消息边界frame机制——一条 protobuf 字节流是流式的、不可自我定界的。因此 Pulsar 协议规定每条消息前都追加一个 4 字节字段指明该帧的大小。单个帧允许的最大尺寸为5 MB。这一点在 Commands.java 中有直接对应的常量// default message size for transfer public static final int DEFAULT_MAX_MESSAGE_SIZE 5 * 1024 * 1024; public static final int MESSAGE_SIZE_FRAME_PADDING 10 * 1024;即 5 MB 的帧上限就是DEFAULT_MAX_MESSAGE_SIZE这也解释了为什么超过默认值的超大消息必须走消息分片chunking机制。Pulsar 协议共定义了两类命令Simple commands简单命令不携带消息负载如Ping、SubscribePayload commands负载命令在发布/投递消息时使用。帧内的 protobuf 命令之后依次是 protobuf 序列化的metadata然后是以原始raw字节形式直接传递的 payload。所有长度字段均为4 字节无符号大端big endian整数。出于效率考虑消息负载采用 raw 格式而非 protobuf 封装——payload 是应用数据的透传无需也不应被协议层重复编码。Simple commands简单命令帧结构组件说明大小字节totalSize帧大小统计其之后的所有字节数4commandSizeprotobuf 序列化后的命令大小4message以 raw 二进制格式存储的 protobuf 消息commandSizePayload commands负载命令帧结构在简单命令的基础上负载命令追加了 magic number、校验和与 metadata组件说明大小字节totalSize帧大小统计其之后的所有字节数4commandSizeprotobuf 序列化后的命令大小4messageprotobuf 消息raw 二进制commandSizemagicNumber2 字节魔数0x0e01标识当前帧格式2checksum其后所有内容的CRC32-C 校验和4metadataSize消息 metadata 的大小4metadata以 protobuf 二进制序列化的消息 metadatametadataSizepayload帧内剩余的全部字节即为 payload可为任意字节序列剩余部分这些结构在仓库源码中得到一一印证。Commands.java 定义了魔数与校验和长度public static final short magicCrc32c 0x0e01; public static final short magicBrokerEntryMetadata 0x0e02; private static final int checksumSize 4;其中magicCrc32c 0x0e01即文档中的魔数表示“后续 4 字节为 CRC32-C 校验和”magicBrokerEntryMetadata 0x0e02则用于携带 Broker 侧补充元数据如broker_timestamp、index对应 proto 中的BrokerEntryMetadata消息的帧。C 客户端 Commands.h 中对帧布局的注释也给出了同样的视图// [TOTAL_SIZE] [CMD_SIZE][CMD] [MAGIC_NUMBER][CHECKSUM] [METADATA_SIZE][METADATA] [PAYLOAD]CRC32-CCastagnoli 变体相比普通 CRC32 拥有更适合硬件加速的查找表实现与“最大传输效率”的设计目标一致。消息 MetadataMessageMetadata消息 metadata 与应用负载一起存储本身是一条 protobuf 序列化消息。metadata 由 Producer 创建并原样传递给 Consumer——Broker 不修改它。字段定义见 PulsarApi.proto 中的MessageMetadata消息。文档中列出的字段如下字段说明producer_name发布该消息的 Producer 名称sequence_idProducer 为消息分配的序列号publish_time发布时戳Unix 时间1970-01-01 UTC 起的毫秒数properties应用自定义的键值对序列使用KeyValue消息对 Pulsar 无特殊含义replicated_from(可选)表明消息经过跨集群复制指明原始发布的集群名称partition_key(可选)发布到分区主题时若存在该 key则以其哈希值决定写入哪个分区compression(可选)标识 payload 被压缩及所使用的压缩算法uncompressed_size(可选)使用压缩时Producer 必须填充压缩前的原始 payload 大小num_messages_in_batch(可选)若该“消息”实为多条消息组成的批次此字段须设为批内消息数对照 PulsarApi.proto 的当前定义可以看到协议在此后版本中还持续演进新增了诸如event_time应用事件时间缺省时可用publish_time代替、encryption_keys/encryption_algo/encryption_param加密、schema_version、ordering_keyKey_Shared 模式下的有序投递键、deliver_at_time定时消息、marker_type内部 marker 消息、事务相关的txnid_least_bits/txnid_most_bits以及消息分片相关的num_chunks_from_msg/total_chunk_msg_size/chunk_id等字段。阅读 proto 是了解协议演进全貌的最可靠途径。compression字段的取值由同文件中的CompressionType枚举给出NONE 0、LZ4 1、ZLIB 2、ZSTD 3、SNAPPY 4。Batch messages批量消息使用批量消息时payload 不再是单条消息体而是一个由若干 entry 组成的列表每个 entry 拥有独立的 metadata由SingleMessageMetadata对象描述定义见 PulsarApi.proto。单个 batch 的 payload 格式如下字段说明metadataSizeN第 N 条单消息 metadata 的 protobuf 序列化大小metadataN第 N 条单消息的 metadataSingleMessageMetadatapayloadN第 N 条消息的原始负载即payload [metadataSize1][metadata1][payload1] [metadataSize2][metadata2][payload2] …逐条重复直到填满帧内剩余空间。每条SingleMessageMetadata的字段字段说明properties应用自定义属性partition key(可选)用于指示哈希到特定分区的 keypayload_size批内该单条消息的 payload 大小启用压缩时整个 batch 被一次性整体压缩而非逐条压缩——这既省去了逐条压缩开销也保证了 batch 内各 entry 的边界信息由未压缩的payload_size保证可解析。交互流程Interactions建立连接客户端向 Broker通常是6650 端口建立 TCP 连接后由客户端负责发起会话。客户端首先发送CommandConnect若 Broker 校验通过了认证则回复Connected连接即可投入使用若认证失败Broker 回复Error命令并直接关闭 TCP 连接。CommandConnect示例message CommandConnect { client_version : Pulsar-Client-Java-v1.15.2, auth_method_name : my-authentication-plugin, auth_data : my-auth-data, protocol_version : 6 }字段说明client_version→ 字符串标识格式不做强制约定auth_method_name→(可选)启用认证时认证插件的名称auth_data→(可选)插件特定的认证数据protocol_version→ 客户端支持的协议版本。Broker 不会下发新协议版本中才引入的命令Broker 也可以强制一个最低协议版本。message CommandConnected { server_version : Pulsar-Broker-v1.15.2, protocol_version : 6 }字段说明server_version→ Broker 版本标识protocol_version→ Broker 支持的协议版本。客户端不得尝试发送比该版本更新引入的命令。这个“双向协商”机制保证了新命令只在双方都支持时才被使用从 Commands.java 的newConnect(...)实现可以看到客户端在建连时就会带上getCurrentProtocolVersion()取ProtocolVersion枚举的最大值以及FeatureFlags如supportsAuthRefresh、supportsBrokerEntryMetadata、supportsPartialProducer把能力声明前置到握手阶段。Keep Alive保活探测为了识别客户端与 Broker 之间长期存在的网络分区或“本机已崩溃但远端 TCP 连接未被中断”的情况断电、内核 panic、强制重启等Pulsar 引入了探测机制双方周期性发送Ping命令如果在超时时间内Broker 侧默认 60 秒未收到Pong响应则关闭 socket。注意实现要求是不对称的一个合格的 Pulsar 客户端不需要主动发送Ping但必须在收到 Broker 的Ping后及时回复Pong以免远端强制断开连接。Commands.java 中提供了现成的newPing()与newPong()构造方法Broker 侧的心跳间隔由 ServiceConfiguration.java 中的keepAliveIntervalSeconds默认 30 秒等参数控制。Producer 交互要发送消息客户端必须先建立 Producer。创建 Producer 时Broker 会先校验该客户端是否**有权authorized**向目标 topic 发布。获得创建成功的确认后客户端即可引用事先协商好的producer_id向 Broker 发布消息。CommandProducermessage CommandProducer { topic : persistent://my-property/my-cluster/my-namespace/my-topic, producer_id : 1, request_id : 1 }参数说明topic→ 要创建 Producer 的完整 topic 名称producer_id→ 客户端生成的 Producer 标识同一连接内需唯一request_id→ 本次请求的标识用于将响应与原始请求配对同一连接内需唯一producer_name→(可选)若指定了 Producer 名称则使用该名称否则由 Broker 生成一个全局唯一的名称。实现上应在 Producer 首次创建时让 Broker 生成名称重连后重建 Producer 时复用该名称。Broker 的响应是ProducerSuccess或Error命令之一。CommandProducerSuccessmessage CommandProducerSuccess { request_id : 1, producer_name : generated-unique-producer-name }参数说明request_id→ 对应CreateProducer请求的原始 idproducer_name→ 生成的全局唯一名称或客户端指定的名称若指定了。CommandSendSend命令用于在已存在的 Producer 上下文中发布一条新消息。它使用“命令 payload”同帧的payload command 帧格式见 payload commands 一节。message CommandSend { producer_id : 1, sequence_id : 0, num_messages : 1 }参数说明producer_id→ 已存在 Producer 的 idsequence_id→ 每条消息都关联一个序列号实现上应为一个从 0 开始的计数器。确认消息实际发布的SendReceipt会通过序列号指代该消息num_messages→(可选)一次性发布一批batch消息时使用。CommandSendReceipt当消息已经在配置的副本数上持久化之后Broker 向 Producer 发送确认回执message CommandSendReceipt { producer_id : 1, sequence_id : 0, message_id : { ledgerId : 123, entryId : 456 } }参数说明producer_id→ 发起发送请求的 Producer idsequence_id→ 已发布消息的序列号message_id→ 系统为已发布消息分配的消息 id在单个集群内唯一。消息 id 由两个 long——ledgerId与entryId——组成这正反映了该唯一 id 是在向 BookKeeper ledger 追加条目时分配的。这一点与 PulsarApi.proto 中的MessageIdData一致除了ledgerId/entryId它还带有partition、batch_index、ack_set等字段分别用于分区主题、批内偏移与批量确认。CommandCloseProducer注意该命令可由 Producer 侧客户端或 Broker 任一方发送。收到CloseProducer时Broker 将停止接收该 Producer 的任何新消息等待所有待处理消息持久化完成后回复Success给客户端。Broker 在优雅故障切换时会主动下发CloseProducer例如Broker 重启或负载均衡器要把 topic 卸载并迁移到其他 Broker。客户端收到该命令后预期会重新执行服务发现lookup并重建 Producer而原 TCP 连接本身不受影响。Consumer 交互Consumer 用于挂载attach到一个订阅subscription并消费其中的消息。每次重连后客户端都需要重新执行订阅如果订阅尚不存在则会被新建。流控Flow controlConsumer 就绪后客户端必须授权give permissionBroker 推送消息通过Flow命令完成。一条Flow命令向 Broker 追加授予推送若干条消息的permit许可。典型的 Consumer 实现使用队列在应用就绪前累积消息当应用消费掉队列中约一半的消息后Consumer 向 Broker 发送与已消费数量相等的 permit 请求补充消息。例如队列大小为 1000消费了 500 条后Consumer 就向 Broker 申请 500 个 permit。这种“许可制”推模型从协议层面防止了慢消费者把 Broker 或网络打爆。Commands.java 中的newFlow(consumerId, messagePermits)即为对应实现。CommandSubscribemessage CommandSubscribe { topic : persistent://my-property/my-cluster/my-namespace/my-topic, subscription : my-subscription-name, subType : Exclusive, consumer_id : 1, request_id : 1 }参数说明topic→ 要创建 Consumer 的完整 topic 名称subscription→ 订阅名称subType→ 订阅类型Exclusive、Shared、Failover、Key_Sharedconsumer_id→ 客户端生成的 Consumer 标识同一连接内需唯一request_id→ 请求标识用于配对请求与响应同一连接内需唯一consumer_name→(可选)客户端可指定 Consumer 名称。该名称可用于在 stats 中追踪特定 Consumer此外在Failover订阅类型中名称用于选举master实际接收消息的那个 ConsumerConsumer 按名称排序第一个被选为 master。CommandFlowmessage CommandFlow { consumer_id : 1, messagePermits : 1000 }参数说明consumer_id→ 已建立 Consumer 的 idmessagePermits→ 授予 Broker 追加推送的消息 permit 数量。CommandMessageMessage命令由 Broker 用于在permit 限额之内向已存在的 Consumer 推送消息。它同样使用携带 payload 的帧格式见 payload commands。message CommandMessage { consumer_id : 1, message_id : { ledgerId : 123, entryId : 456 } }CommandAckAck用于向 Broker 指示某条消息已被应用成功处理、可以被丢弃。同时Broker 会基于已确认的消息维护消费位置consumer position。message CommandAck { consumer_id : 1, ack_type : Individual, message_id : { ledgerId : 123, entryId : 456 } }参数说明consumer_id→ 已建立 Consumer 的 idack_type→ 确认类型Individual逐条确认或Cumulative累积确认message_id→ 要确认的消息 idvalidation_error→(可选)表明消费者因如下原因丢弃了消息UncompressedSizeCorruption、DecompressionError、ChecksumMismatch、BatchDeSerializeError。CommandCloseConsumer注意该命令可由客户端或 Broker 任一方发送行为与CloseProducer完全相同。CommandRedeliverUnacknowledgedMessagesConsumer 可以要求 Broker 重投redeliver已被推送但尚未确认的部分或全部待处理消息。该 protobuf 对象接受一个消息 id 列表若列表为空Broker 将重投全部待处理消息。重投时消息可以发给同一个 Consumer也可以在 Shared 订阅场景下分散到所有可用 Consumer 上。CommandReachedEndOfTopic当 topic 已被“终止terminated”且该订阅上的所有消息都已被确认时Broker 向特定 Consumer 发送此命令。客户端应据此通知应用“不会再有消息从该 Consumer 到达”。CommandConsumerStats该命令由客户端发送用于从 Broker 拉取Subscriber 与 Consumer 级别的统计信息。参数request_id→ 请求 id用于关联请求与响应consumer_id→ 已建立 Consumer 的 id。CommandConsumerStatsResponseBroker 对ConsumerStats请求的响应内容是对应consumer_id的 Subscriber 与 Consumer 级别统计。若设置了error_code或error_message字段则表明请求失败。CommandUnsubscribe该命令由客户端发送将consumer_id从关联的 topic 上退订。参数request_id→ 请求 idconsumer_id→ 需要退订的已建立 Consumer 的 id。服务发现Service DiscoveryTopic Lookup每当客户端需要创建或重连一个 Producer / Consumer 时都要先执行一次Topic Lookup用于发现当前是哪个 Broker 正在服务目标 topic。Lookup 也可以通过管理 APIREST完成参见 admin API 文档自 Pulsar 1.16 起Lookup 也可以直接在二进制协议内部完成。以如下部署为例服务发现组件运行在pulsar://broker.example.com:6650各个具体 Broker 运行在pulsar://broker-1.example.com:6650、pulsar://broker-2.example.com:6650……客户端连接到发现地址后下发LookupTopic命令响应要么是“应该连接的 Broker 地址”要么是“应该重试 lookup 的地址”重定向。LookupTopic必须在使用过Connect/Connected初始握手的连接上发送。message CommandLookupTopic { topic : persistent://my-property/my-cluster/my-namespace/my-topic, request_id : 1, authoritative : false }字段说明topic→ 要查询的 topic 名称request_id→ 请求 id会随响应原样返回authoritative→ 首次 lookup 请求应使用false跟随重定向响应时客户端应传入响应中所含的同一值。LookupTopicResponse的成功响应示例message CommandLookupTopicResponse { request_id : 1, response : Connect, brokerServiceUrl : pulsar://broker-1.example.com:6650, brokerServiceUrlTls : pulsarssl://broker-1.example.com:6651, authoritative : true }重定向Redirect响应示例message CommandLookupTopicResponse { request_id : 1, response : Redirect, brokerServiceUrl : pulsar://broker-2.example.com:6650, brokerServiceUrlTls : pulsarssl://broker-2.example.com:6651, authoritative : true }第二种情况下客户端需要向broker-2.example.com重新发起LookupTopic请求——该 Broker 将能给出确定性的答案。这种“redirect 直至 authoritative”的机制使得元数据不一致期间例如 topic 正在 Broker 间迁移lookup 仍能收敛到正确节点。Commands.java 中的newLookup(topic, authoritative, requestId)与newLookupResponse(...)正是客户端/服务端两侧构造这些命令的入口响应中response字段即 proto 里的LookupTypeConnect/Redirect。Partitioned topics discovery分区主题发现分区主题元数据发现用于确定某 topic 是否为“partitioned topic”以及配置了多少个分区。若 topic 被标记为分区主题客户端需要为每个分区各创建一个 Producer 或 Consumertopic 名称使用partition-X后缀。该信息只在首次创建 Producer / Consumer 时需要获取重连后无需重复。其工作机制与 topic lookup 类似客户端向服务发现地址发送请求响应中携带实际的元数据。CommandPartitionedTopicMetadatamessage CommandPartitionedTopicMetadata { topic : persistent://my-property/my-cluster/my-namespace/my-topic, request_id : 1 }字段说明topic→ 要检查分区元数据的 topicrequest_id→ 请求 id会随响应原样返回。CommandPartitionedTopicMetadataResponse携带元数据的响应示例message CommandPartitionedTopicMetadataResponse { request_id : 1, response : Success, partitions : 32 }Protobuf 接口与延伸阅读Pulsar 的全部 Protobuf 定义都集中在 PulsarApi.proto 这一个文件中包名pulsar.protooptimize_for LITE_RUNTIME面向轻量运行时优化。文中出现的BaseCommand、MessageMetadata、SingleMessageMetadata、MessageIdData、CompressionType、KeySharedMode等类型均可在其中检索到完整定义C 客户端的对应实现在 pulsar-client-cpp/lib/Commands.h 与 pulsar-client-cpp/lib/Commands.cc 中可作为跨语言实现交叉验证帧格式与命令语义的参照。理解这套二进制协议的价值在于它解释了 Pulsar 客户端行为背后的每一条“规则”——为什么重连后要重新 lookup 与 subscribeLookupTopic/Subscribe语义、为什么有redeliverUnacknowledgedMessagesAPIRedeliverUnacknowledgedMessages命令、为什么批量压缩是整批一起压MessageMetadata.compressionnum_messages_in_batch语义、为什么消息 id 是(ledgerId, entryId)二元组BookKeeper 持久化模型。掌握这些协议细节后无论是排查“消息重复投递”“permit 饥饿”“心跳断连”一类问题还是自研网关/代理都有了可靠的协议级依据。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐React Native Push Notification iOS未来展望新功能和技术趋势分析React Native Push Notification iOS未来展望新功能和技术趋势分析 React Native Push NotificationApache Pulsar 二进制协议规范深度解析帧格式、命令交互与服务发现Apache Pulsar 二进制协议规范深度解析帧格式、命令交互与服务发现 Pulsar 的生产者/消费者与 Broker 之间通过一套自定义的二进制协议进消息队列后端流处理Apache Pulsar 二进制协议Binary Protocol深度解析从帧格式到命令交互全指南Apache Pulsar 二进制协议Binary Protocol深度解析从帧格式到命令交互全指南 output_article Apache Pul消息队列后端流处理上一篇NodeJS JWT Authentication Sample与Postman集成高效测试API的实用技巧下一篇终极指南FingerprintJS 浏览器指纹库免费使用教程创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

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

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

免费获取报价 →
↑