news 2026/9/25 2:46:01

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

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Pulsar 二进制协议规范详解:从帧结构到 Lookup 的服务端通信原理
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache 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 的 BaseCommand

Pulsar 客户端和 Broker 之间交换的基本单元是命令(command)。所有命令都被封装为二进制的Protocol Buffers(protobuf)消息,其格式由仓库中的 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 协议共定义了两类命令:

  1. Simple commands(简单命令):不携带消息负载,如Ping、Subscribe;
  2. Payload commands(负载命令):在发布/投递消息时使用。帧内的 protobuf 命令之后依次是 protobuf 序列化的metadata,然后是以原始(raw)字节形式直接传递的 payload。所有长度字段均为4 字节无符号大端(big endian)整数。

出于效率考虑,消息负载采用 raw 格式而非 protobuf 封装——payload 是应用数据的透传,无需也不应被协议层重复编码。

Simple commands(简单命令)帧结构

组件说明大小(字节)
totalSize帧大小,统计其之后的所有字节数4
commandSizeprotobuf 序列化后的命令大小4
message以 raw 二进制格式存储的 protobuf 消息commandSize

Payload commands(负载命令)帧结构

在简单命令的基础上,负载命令追加了 magic number、校验和与 metadata:

组件说明大小(字节)
totalSize帧大小,统计其之后的所有字节数4
commandSizeprotobuf 序列化后的命令大小4
messageprotobuf 消息(raw 二进制)commandSize
magicNumber2 字节魔数0x0e01,标识当前帧格式2
checksum其后所有内容的CRC32-C 校验和4
metadataSize消息 metadata 的大小4
metadata以 protobuf 二进制序列化的消息 metadatametadataSize
payload帧内剩余的全部字节即为 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-C(Castagnoli 变体)相比普通 CRC32 拥有更适合硬件加速的查找表实现,与“最大传输效率”的设计目标一致。

消息 Metadata(MessageMetadata)

消息 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_key(Key_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 条单消息的 metadata(SingleMessageMetadata)
payloadN第 N 条消息的原始负载

即payload = [metadataSize1][metadata1][payload1] [metadataSize2][metadata2][payload2] …,逐条重复直到填满帧内剩余空间。

每条SingleMessageMetadata的字段:

字段说明
properties应用自定义属性
partition key(可选)用于指示哈希到特定分区的 key
payload_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 发布消息。

CommandProducer
message 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命令之一。

CommandProducerSuccess
message CommandProducerSuccess { "request_id" : 1, "producer_name" : "generated-unique-producer-name" }

参数说明:

  • request_id→ 对应CreateProducer请求的原始 id;
  • producer_name→ 生成的全局唯一名称,或客户端指定的名称(若指定了)。
CommandSend

Send命令用于在已存在的 Producer 上下文中发布一条新消息。它使用“命令 + payload”同帧的payload command 帧格式(见 payload commands 一节)。

message CommandSend { "producer_id" : 1, "sequence_id" : 0, "num_messages" : 1 }

参数说明:

  • producer_id→ 已存在 Producer 的 id;
  • sequence_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 id;
  • sequence_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 control)

Consumer 就绪后,客户端必须授权(give permission)Broker 推送消息,通过Flow命令完成。一条Flow命令向 Broker 追加授予推送若干条消息的permit(许可)。典型的 Consumer 实现使用队列在应用就绪前累积消息;当应用消费掉队列中约一半的消息后,Consumer 向 Broker 发送与已消费数量相等的 permit 请求补充消息。例如队列大小为 1000,消费了 500 条后,Consumer 就向 Broker 申请 500 个 permit。

这种“许可制”推模型从协议层面防止了慢消费者把 Broker 或网络打爆。Commands.java 中的newFlow(consumerId, messagePermits)即为对应实现。

CommandSubscribe
message 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_Shared;
  • consumer_id→ 客户端生成的 Consumer 标识,同一连接内需唯一;
  • request_id→ 请求标识,用于配对请求与响应,同一连接内需唯一;
  • consumer_name→(可选)客户端可指定 Consumer 名称。该名称可用于在 stats 中追踪特定 Consumer;此外,在Failover订阅类型中,名称用于选举master(实际接收消息的那个 Consumer):Consumer 按名称排序,第一个被选为 master。
CommandFlow
message CommandFlow { "consumer_id" : 1, "messagePermits" : 1000 }

参数说明:

  • consumer_id→ 已建立 Consumer 的 id;
  • messagePermits→ 授予 Broker 追加推送的消息 permit 数量。
CommandMessage

Message命令由 Broker 用于在permit 限额之内向已存在的 Consumer 推送消息。它同样使用携带 payload 的帧格式(见 payload commands)。

message CommandMessage { "consumer_id" : 1, "message_id" : { "ledgerId" : 123, "entryId" : 456 } }
CommandAck

Ack用于向 Broker 指示某条消息已被应用成功处理、可以被丢弃。同时,Broker 会基于已确认的消息维护消费位置(consumer position)。

message CommandAck { "consumer_id" : 1, "ack_type" : "Individual", "message_id" : { "ledgerId" : 123, "entryId" : 456 } }

参数说明:

  • consumer_id→ 已建立 Consumer 的 id;
  • ack_type→ 确认类型:Individual(逐条确认)或Cumulative(累积确认);
  • message_id→ 要确认的消息 id;
  • validation_error→(可选)表明消费者因如下原因丢弃了消息:UncompressedSizeCorruption、DecompressionError、ChecksumMismatch、BatchDeSerializeError。
CommandCloseConsumer

注意:该命令可由客户端或 Broker 任一方发送,行为与CloseProducer完全相同。

CommandRedeliverUnacknowledgedMessages

Consumer 可以要求 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。
CommandConsumerStatsResponse

Broker 对ConsumerStats请求的响应,内容是对应consumer_id的 Subscriber 与 Consumer 级别统计。若设置了error_code或error_message字段,则表明请求失败。

CommandUnsubscribe

该命令由客户端发送,将consumer_id从关联的 topic 上退订。参数:

  • request_id→ 请求 id;
  • consumer_id→ 需要退订的已建立 Consumer 的 id。

服务发现(Service Discovery)

Topic Lookup

每当客户端需要创建或重连一个 Producer / Consumer 时,都要先执行一次Topic Lookup,用于发现当前是哪个 Broker 正在服务目标 topic。Lookup 也可以通过管理 API(REST)完成(参见 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" : "pulsar+ssl://broker-1.example.com:6651", "authoritative" : true }

重定向(Redirect)响应示例:

message CommandLookupTopicResponse { "request_id" : 1, "response" : "Redirect", "brokerServiceUrl" : "pulsar://broker-2.example.com:6650", "brokerServiceUrlTls" : "pulsar+ssl://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 里的LookupType(Connect/Redirect)。

Partitioned topics discovery(分区主题发现)

分区主题元数据发现用于确定某 topic 是否为“partitioned topic”以及配置了多少个分区。若 topic 被标记为分区主题,客户端需要为每个分区各创建一个 Producer 或 Consumer,topic 名称使用partition-X后缀。该信息只在首次创建 Producer / Consumer 时需要获取,重连后无需重复。

其工作机制与 topic lookup 类似:客户端向服务发现地址发送请求,响应中携带实际的元数据。

CommandPartitionedTopicMetadata
message CommandPartitionedTopicMetadata { "topic" : "persistent://my-property/my-cluster/my-namespace/my-topic", "request_id" : 1 }

字段说明:

  • topic→ 要检查分区元数据的 topic;
  • request_id→ 请求 id,会随响应原样返回。
CommandPartitionedTopicMetadataResponse

携带元数据的响应示例:

message CommandPartitionedTopicMetadataResponse { "request_id" : 1, "response" : "Success", "partitions" : 32 }

Protobuf 接口与延伸阅读

Pulsar 的全部 Protobuf 定义都集中在 PulsarApi.proto 这一个文件中(包名pulsar.proto,optimize_for = LITE_RUNTIME,面向轻量运行时优化)。文中出现的BaseCommand、MessageMetadata、SingleMessageMetadata、MessageIdData、CompressionType、KeySharedMode等类型均可在其中检索到完整定义;C++ 客户端的对应实现在 pulsar-client-cpp/lib/Commands.h 与 pulsar-client-cpp/lib/Commands.cc 中,可作为跨语言实现交叉验证帧格式与命令语义的参照。

理解这套二进制协议的价值在于:它解释了 Pulsar 客户端行为背后的每一条“规则”——为什么重连后要重新 lookup 与 subscribe(LookupTopic/Subscribe语义)、为什么有redeliverUnacknowledgedMessagesAPI(RedeliverUnacknowledgedMessages命令)、为什么批量压缩是整批一起压(MessageMetadata.compression+num_messages_in_batch语义)、为什么消息 id 是(ledgerId, entryId)二元组(BookKeeper 持久化模型)。掌握这些协议细节后,无论是排查“消息重复投递”“permit 饥饿”“心跳断连”一类问题,还是自研网关/代理,都有了可靠的协议级依据。

  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载
上一篇:NodeJS JWT Authentication Sample与Postman集成:高效测试API的实用技巧
下一篇:终极指南:FingerprintJS 浏览器指纹库免费使用教程

创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考

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

jc 解析 efibootmgr:将 Linux UEFI 启动项管理输出转换为 JSON

开发工具 【免费下载链接】jc CLI tool and python library that converts the output of popular command-line tools, file-types, and common strings to JSON, YAML, or Dictionaries. This allows piping of output to tools like jq and simplifying automation scripts.…

作者头像 李华
网站建设 2026/9/25 2:44:59

澎湃OS时代BL锁机制深度解析与绕过实践

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

作者头像 李华
网站建设 2026/9/25 2:44:14

MoneyPrinterTurbo:输入一个主题,3 步 AI 自动出片

MoneyPrinterTurbo:输入一个主题,3 步 AI 自动出片 【免费下载链接】MoneyPrinterTurbo 利用 AI 大模型和自动化工作流,根据主题或关键词一键生成高清短视频。Generate HD short videos from a topic or keyword with an automated AI workfl…

作者头像 李华