- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
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 协议共定义了两类命令:
- Simple commands(简单命令):不携带消息负载,如
Ping、Subscribe; - Payload commands(负载命令):在发布/投递消息时使用。帧内的 protobuf 命令之后依次是 protobuf 序列化的metadata,然后是以原始(raw)字节形式直接传递的 payload。所有长度字段均为4 字节无符号大端(big endian)整数。
出于效率考虑,消息负载采用 raw 格式而非 protobuf 封装——payload 是应用数据的透传,无需也不应被协议层重复编码。
Simple commands(简单命令)帧结构
| 组件 | 说明 | 大小(字节) |
|---|---|---|
totalSize | 帧大小,统计其之后的所有字节数 | 4 |
commandSize | protobuf 序列化后的命令大小 | 4 |
message | 以 raw 二进制格式存储的 protobuf 消息 | commandSize |
Payload commands(负载命令)帧结构
在简单命令的基础上,负载命令追加了 magic number、校验和与 metadata:
| 组件 | 说明 | 大小(字节) |
|---|---|---|
totalSize | 帧大小,统计其之后的所有字节数 | 4 |
commandSize | protobuf 序列化后的命令大小 | 4 |
message | protobuf 消息(raw 二进制) | commandSize |
magicNumber | 2 字节魔数0x0e01,标识当前帧格式 | 2 |
checksum | 其后所有内容的CRC32-C 校验和 | 4 |
metadataSize | 消息 metadata 的大小 | 4 |
metadata | 以 protobuf 二进制序列化的消息 metadata | metadataSize |
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_id | Producer 为消息分配的序列号 |
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
相关推荐
React Native Push Notification iOS未来展望:新功能和技术趋势分析
React Native Push Notification iOS未来展望:新功能和技术趋势分析 React Native Push Notification
Apache Pulsar 二进制协议规范深度解析:帧格式、命令交互与服务发现
Apache Pulsar 二进制协议规范深度解析:帧格式、命令交互与服务发现 Pulsar 的生产者/消费者与 Broker 之间通过一套自定义的二进制协议进
消息队列后端流处理Apache Pulsar 二进制协议(Binary Protocol)深度解析:从帧格式到命令交互全指南
Apache Pulsar 二进制协议(Binary Protocol)深度解析:从帧格式到命令交互全指南 <output_article Apache Pul
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考