- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
Apache Pulsar 的所有消息数据最终都落盘在 Apache BookKeeper 的 ledger(分段日志)中。为了让运维人员与开发者能快速识别某条 ledger 属于哪个 topic、哪个 cursor、是否承载 schema 或消息压缩(compaction)产物,Pulsar 在创建每条 ledger 时会附加一组自定义元数据(custom metadata),并持久化在 ZooKeeper 上。本文以version-2.3.1官方 cookbook《BookKeeper Ledger Metadata》为骨架,结合当前仓库中的源码实现,逐字段解释这些元数据的含义、写入位置与读取方式,帮助你在排障、数据迁移和存储审计时快速定位"这条 ledger 到底存的是什么"。
Pulsar 为什么需要 ledger 元数据
Pulsar 将 topic 的消息以分段(segment)形式写入 BookKeeper ledger:一个 managed-ledger 对应一个或多个 ledger,topic 的每个游标(cursor)也会单独开辟 ledger 记录消费位点,topic 压缩(compaction)会生成新的压缩后 ledger,schema 注册信息同样存储在独立 ledger 中。当这些 ledger 混存在同一套 BookKeeper 集群里时,仅靠 ledger id 无法判断其用途。
Pulsar 的解法是在创建 ledger 时附加自定义元数据(custom metadata),这些元数据与 ledger 自身属性(如 ensemble size、写入/确认 quorum)一起写入 BookKeeper 的 ledger 元数据节点,最终持久化在 ZooKeeper 上,并可通过 BookKeeper API 读取。官方 cookbook 明确指出:
Pulsar stores data on BookKeeper ledgers, you can understand the contents of a ledger by inspecting the metadata attached to the ledger. Such metadata are stored on ZooKeeper and they are readable using BookKeeper APIs.
当前元数据字段总览
原文档给出的元数据描述如下:
| 作用范围(Scope) | 元数据名(Metadata name) | 元数据值(Metadata value) |
|---|---|---|
| 所有 ledger | application | 'pulsar' |
| 所有 ledger | component | 'managed-ledger'、'schema'、'compacted-topic' |
| Managed ledgers | pulsar/managed-ledger | ledger 的名称(name of the ledger) |
| Cursor | pulsar/cursor | cursor 的名称(name of the cursor) |
| Compacted topic | pulsar/compactedTopic | 原始 topic 的名称(name of the original topic) |
| Compacted topic | pulsar/compactedTo | 最后一条已压缩消息的 id(id of the last compacted message) |
对照当前仓库源码(见下文LedgerMetadataUtils),可以确认这套字段体系被完整保留,且新增了pulsar/schemaId用于 schema ledger。需要特别指出一个版本差异:version-2.3.1文档中 compacted ledger 的component值写作'compacted-topic',而当前仓库源码中的实际取值为'compacted-ledger',阅读旧文档或对接旧版本集群时需要注意这一命名差异。
逐项解读:每个元数据字段的含义与来源
application:标记元数据归属应用
所有 Pulsar 创建的 ledger 都会带上application = 'pulsar',用于在混合部署场景中区分哪些 ledger 由 Pulsar 管理、哪些来自其他 BookKeeper 用户(例如直接使用 BookKeeper 的作业)。它是 Pulsar 写入元数据时的"签名"。
component:标记 ledger 的业务组件类型
component标识这条 ledger 承担的业务角色,当前源码中支持三种取值:
managed-ledger:普通消息数据 ledger;compacted-ledger:topic 压缩(compaction)后生成的 ledger(旧版文档记为compacted-topic);schema:存储 schema 注册信息的 ledger。
pulsar/managed-ledger:归属的 managed-ledger 名称
该字段标记数据 ledger 属于哪个 managed-ledger。Pulsar 的 managed-ledger 名称通常是tenant/namespace/topic的持久化名称,因此通过该字段可以直接推断这条 ledger 承载的是哪个 topic 的消息数据。
pulsar/cursor:归属的 cursor 名称
cursor 是 Pulsar 消费进度的逻辑游标(对应一个订阅或 reader)。每个 cursor 会额外创建 ledger 来持久化其消费位点,pulsar/cursor记录该 ledger 对应的 cursor 名称。通过该字段可以判断一条 ledger 是游标位点 ledger 而非消息数据 ledger。
pulsar/compactedTopic与pulsar/compactedTo:压缩产物的原始 topic 与位点
启用 topic 压缩后,Pulsar 的TwoPhaseCompactor会把每个 key 的最新消息写入一条全新的 ledger,并在其上附加:
pulsar/compactedTopic:压缩前原始 topic 的名称;pulsar/compactedTo:压缩后最后一条消息的 MessageId(以二进制字节形式存储)。
有了这两个字段,压缩产物 ledger 就能精确回溯"它来自哪个 topic、压缩截止到哪条消息"。
源码级实现:元数据在何处被写入
当前仓库中,所有 ledger 元数据的构建逻辑都集中在 LedgerMetadataUtils.java 这一个工具类中,它以Map<String, byte[]>(key 为字符串、value 为 UTF-8/二进制字节)的形式产出元数据,并原样传给 BookKeeper 的createLedgerAPI。类内定义的常量与文档表格一一对应(见 LedgerMetadataUtils.java)。
基础元数据:buildBaseManagedLedgerMetadata
static Map<String, byte[]> buildBaseManagedLedgerMetadata(String name) { return ImmutableMap.of( METADATA_PROPERTY_APPLICATION, METADATA_PROPERTY_APPLICATION_PULSAR, METADATA_PROPERTY_COMPONENT, METADATA_PROPERTY_COMPONENT_MANAGED_LEDGER, METADATA_PROPERTY_MANAGED_LEDGER_NAME, name.getBytes(StandardCharsets.UTF_8)); }即每个 managed-ledger 打开新 ledger 时都会写入application、component、pulsar/managed-ledger三对键值。调用点在 ManagedLedgerImpl.java:this.ledgerMetadata = LedgerMetadataUtils.buildBaseManagedLedgerMetadata(name);,之后该 Map 会作为默认自定义元数据随每次asyncCreateLedger一并提交。
cursor 元数据:buildAdditionalMetadataForCursor
static Map<String, byte[]> buildAdditionalMetadataForCursor(String name) { return ImmutableMap.of(METADATA_PROPERTY_CURSOR_NAME, name.getBytes(StandardCharsets.UTF_8)); }只附加pulsar/cursor一个字段,与基础元数据合并后写入 cursor 的位点 ledger。调用点在 ManagedCursorImpl.java,创建新 ledger 时通过LedgerMetadataUtils.buildAdditionalMetadataForCursor(name)传入 cursor 名称。
compacted ledger 元数据:buildMetadataForCompactedLedger
public static Map<String, byte[]> buildMetadataForCompactedLedger(String compactedTopic, byte[] compactedToMessageId) { return ImmutableMap.of( METADATA_PROPERTY_APPLICATION, METADATA_PROPERTY_APPLICATION_PULSAR, METADATA_PROPERTY_COMPONENT, METADATA_PROPERTY_COMPONENT_COMPACTED_LEDGER, METADATA_PROPERTY_COMPACTEDTOPIC, compactedTopic.getBytes(StandardCharsets.UTF_8), METADATA_PROPERTY_COMPACTEDTO, compactedToMessageId ); }注意pulsar/compactedTo的 value 是 MessageId 的字节数组(to.toByteArray()),并非可读字符串。调用点在 TwoPhaseCompactor.java 的phaseTwo阶段——两阶段压缩器完成"为每个 key 保留最新消息"的扫描后,即以此元数据创建压缩产物 ledger,随后把整理后的消息逐条写入。
schema 元数据:buildMetadataForSchema(当前版本新增)
public static Map<String, byte[]> buildMetadataForSchema(String schemaId) { return ImmutableMap.of( METADATA_PROPERTY_APPLICATION, METADATA_PROPERTY_APPLICATION_PULSAR, METADATA_PROPERTY_COMPONENT, METADATA_PROPERTY_COMPONENT_SCHEMA, METADATA_PROPERTY_SCHEMAID, schemaId.getBytes(StandardCharsets.UTF_8) ); }用于 schema 存储 ledger,额外携带pulsar/schemaId(该字段未出现在 2.3.1 文档表格中,是后续版本补充的能力)。调用点在 BookkeeperSchemaStorage.java:创建 schema ledger 时构造metadataMap,并通过bookKeeper.asyncCreateLedger(..., null, metadata)的最后一个参数传给 BookKeeper。
一个补充:placement policy 配置元数据
同一工具类还提供buildMetadataForPlacementPolicyConfig,把EnsemblePlacementPolicy的实现类与属性编码后写入元数据(见 LedgerMetadataUtils.java),由 ManagedLedgerImpl.java 在覆盖默认放置策略时调用。这类元数据用于让 BookKeeper 在后续写入时沿用相同的副本放置策略。
如何读取这些元数据
按照官方 cookbook 的说明,元数据存储在 ZooKeeper 上,并且可以通过 BookKeeper API 读取。具体落地方式分两层:
存储层(ZooKeeper):每条 ledger 的元数据(包括自定义元数据)以 znode 形式保存在 BookKeeper 的 ledger 元数据目录(默认
/ledgers)下,由 BookKeeper 的LedgerManager负责读写,Pulsar 自身不直接访问该目录。读取层(BookKeeper API):元数据是创建 ledger 时通过
createLedger(..., properties)传入的键值对集合。读取侧可借助 BookKeeper 提供的 Admin API / CLI 获取LedgerMetadata对象,其自定义元数据部分(custom metadata)即上述各字段。Pulsar 侧的实际写路径可从 BookkeeperSchemaStorage.java 看到:asyncCreateLedger的最后一个参数就是承载这些字段的元数据 Map。
说明:元数据值多为 UTF-8 字符串,唯有
pulsar/compactedTo是 MessageId 的二进制序列化结果,直接查看 ZooKeeper 或 CLI 输出时看到的是乱码字节,需用 Pulsar 的 MessageId 解析逻辑还原。
实战应用:用元数据快速定位问题
理解这套元数据后,可以在以下场景直接受益:
- 排查磁盘占用:按
component区分消息数据、压缩产物与 schema 数据 ledger,定位哪类数据占用了大量存储; - 定位 topic 数据:根据
pulsar/managed-ledger的名称前缀反查具体tenant/namespace/topic,确认某条 ledger 是否属于待迁移或待清理的 topic; - 校验压缩结果:对比压缩产物 ledger 上的
pulsar/compactedTopic与pulsar/compactedTo,确认 compaction 是否按预期覆盖了目标 topic 与截止位点; - 区分游标与数据:识别
pulsar/cursor类型的 ledger,避免在数据清理时误删消费位点导致重复消费。
同时需要注意版本差异:version-2.3.1文档中 compacted ledger 的component值为'compacted-topic',当前仓库源码中的实际值为'compacted-ledger',且 schema ledger 增加了pulsar/schemaId字段。以源码 LedgerMetadataUtils.java 为准,可以保证你的解析逻辑与当前版本 Pulsar 完全一致。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar BookKeeper Ledger 元数据完全指南:如何通过 ZooKeeper 与 BookKeeper API 解读数据存储结构
Apache Pulsar BookKeeper Ledger 元数据完全指南:如何通过 ZooKeeper 与 BookKeeper API 解读数据存储结构
消息队列后端流处理Apache Pulsar BookKeeper Ledger 元数据解析:字段含义、存储位置与源码实现
Apache Pulsar BookKeeper Ledger 元数据解析:字段含义、存储位置与源码实现 本文以 Apache Pulsar 的 BookKee
消息队列后端流处理Apache Pulsar 架构总览:Broker、BookKeeper、元数据存储与服务发现全解析
Apache Pulsar 架构总览:Broker、BookKeeper、元数据存储与服务发现全解析 本文基于 Apache Pulsar 官方架构文档( si
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考