☰
Apache Pulsar BookKeeper Ledger 元数据解析:从 ZooKeeper 到源码的完整指南
2026/9/26 19:36:20 网站建设 项目流程
  • 消息队列
  • 后端
  • 流处理

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

本指南以 Apache Pulsar 官方 Cookbook《BookKeeper Ledger Metadata》为主体,深入讲解 Pulsar 如何通过 BookKeeper ledger 自定义元数据(metadata)标记每一份数据卷的用途与归属。读完本文,你将掌握 ledger 元数据在 ZooKeeper 中的存储形态、每一类元数据字段的语义与取值,以及如何通过 BookKeeper API 读取这些元数据来排查与运维 Pulsar 的存储层。

一、背景:为什么需要查看 Ledger 元数据

Pulsar 的所有消息数据最终都落盘在 BookKeeper 的 ledger 上。一个运行中的 Pulsar 集群会同时维护大量 ledger,它们各自承担不同职责:有的是主题(topic)数据的 managed ledger,有的是消费游标(cursor)持久化状态,有的是 topic 压缩(compaction)产生的专用 ledger,还有的是 Schema 存储 ledger。

仅凭一个数值型的 ledger id,运维人员很难判断"这个 ledger 到底存的是什么"。为此,Pulsar 在创建每一个 ledger 时都会附加一组自定义元数据(custom metadata),将 ledger 的应用归属、组件类型、对应实体名称等信息写入其中。这些元数据保存在 ZooKeeper 上,并且可以通过 BookKeeper 的标准 API 读取出来——这正是排查数据、分析存储分布、诊断压缩或游标问题时的重要入口。

说明:本指南对应的官方文档版本为 2.3.0(见 cookbooks-bookkeepermetadata.md),而当前仓库主干版本为 2.10.6-SNAPSHOT(见 pom.xml)。文中会同时给出两版信息,并以当前源码为准标注差异。

二、元数据存在哪里:ZooKeeper + BookKeeper API

原文档明确指出:

  • Pulsar 将数据存储在 BookKeeper ledgers 上,你可以通过检查 ledger 附加的元数据来理解该 ledger 的内容;
  • 这些元数据存储在 ZooKeeper 上;
  • 它们可以使用 BookKeeper API 读取。

在 BookKeeper 的存储模型中,每个 ledger 的元数据(包括创建时间、ensemble、写入 quorum、ack quorum、digest 类型以及自定义属性)都会持久化到 ZooKeeper 的元数据节点中。Pulsar 在调用asyncCreateLedger时,会通过metadata参数把自定义属性一并写入,随后这些属性便与 ledger 本体绑定,可通过LedgerHandle.getLedgerMetadata()获取。

从当前源码看,Pulsar 暴露了一条便捷的查询链路:ManagedLedgerImpl.getLedgerMetadata(long ledgerId)(见 ManagedLedgerImpl.java)——对于当前活跃的 ledger 直接返回currentLedger.getLedgerMetadata().toSafeString(),对于已滚动的旧 ledger 则通过getLedgerHandle(ledgerId)打开后读取其元数据。该查询结果也会出现在getManagedLedgerInternalStats(boolean includeLedgerMetadata)的管理统计中(见 ManagedLedgerImpl.java),便于管理员通过 Pulsar Admin 接口直接审视每个 ledger 的元数据内容。

三、当前元数据字段总览

原文档给出了当时(2.3.0)全部元数据字段的权威说明,下表完整继承并补充了取值说明:

作用域(Scope)元数据名(Metadata name)元数据值(Metadata value)
所有 ledgerapplication'pulsar'
所有 ledgercomponent'managed-ledger'、'schema'、'compacted-topic'
Managed ledgerspulsar/managed-ledgerledger 的名称(name of the ledger)
Cursorpulsar/cursor游标名称(name of the cursor)
Compacted topicpulsar/compactedTopic原始主题名称(name of the original topic)
Compacted topicpulsar/compactedTo最后一条已压缩消息的 id(id of the last compacted message)

其中:

  • application与component是每个 Pulsar 创建的 ledger 都会携带的基础字段,用于标识"这是 Pulsar 写的数据、属于哪个组件";
  • pulsar/managed-ledger标记该 ledger 属于哪一个 managed ledger(即哪个 topic);
  • pulsar/cursor标记游标持久化 ledger 对应的游标名;
  • pulsar/compactedTopic与pulsar/compactedTo仅出现在压缩(compaction)产生的 ledger 上,分别记录被压缩的原始主题和压缩完成后最后一条消息的位置。

四、源码级实现:LedgerMetadataUtils

当前仓库中,所有这些元数据键名与构造逻辑都集中在 LedgerMetadataUtils.java 内。该类是"Utilities for managing BookKeeper Ledgers custom metadata",定义了全部键常量:

常量键名用途
METADATA_PROPERTY_APPLICATIONapplication应用标识,固定为pulsar
METADATA_PROPERTY_COMPONENTcomponent组件标识
METADATA_PROPERTY_MANAGED_LEDGER_NAMEpulsar/managed-ledgermanaged ledger 名称
METADATA_PROPERTY_CURSOR_NAMEpulsar/cursor游标名称
METADATA_PROPERTY_COMPACTEDTOPICpulsar/compactedTopic被压缩的原始主题
METADATA_PROPERTY_COMPACTEDTOpulsar/compactedTo最后一条已压缩消息 id
METADATA_PROPERTY_SCHEMAIDpulsar/schemaIdSchema id(新增于 2.3.0 之后的版本)

版本差异提示:原文档将component的取值列举为'managed-ledger'、'schema'、'compacted-topic';而在当前源码(2.10.6-SNAPSHOT)中,压缩 ledger 的组件值实际写作'compacted-ledger'(见 LedgerMetadataUtils.java),并且新增了pulsar/schemaId这一键(见 LedgerMetadataUtils.java)。排查时请以实际部署版本的取值为准。

4.1 各类 ledger 的元数据构造

LedgerMetadataUtils提供了四个核心构造方法,分别对应原文档表格中的各个作用域:

1. Managed ledger 基础元数据

static Map<String, byte[]> buildBaseManagedLedgerMetadata(String name) { return ImmutableMap.of( METADATA_PROPERTY_APPLICATION, METADATA_PROPERTY_APPLICATION_PULSAR, // application=pulsar METADATA_PROPERTY_COMPONENT, METADATA_PROPERTY_COMPONENT_MANAGED_LEDGER, // component=managed-ledger METADATA_PROPERTY_MANAGED_LEDGER_NAME, name.getBytes(StandardCharsets.UTF_8)); // pulsar/managed-ledger=<name> }

该方法在ManagedLedgerImpl构造时即被调用,成为该 managed ledger 一切 ledger 的默认元数据(见 ManagedLedgerImpl.java),因此每个 managed ledger 下的数据 ledger 都天然携带application=pulsar、component=managed-ledger、pulsar/managed-ledger=<ledger名>三组键值。

2. Cursor 附加元数据

static Map<String, byte[]> buildAdditionalMetadataForCursor(String name) { return ImmutableMap.of(METADATA_PROPERTY_CURSOR_NAME, name.getBytes(StandardCharsets.UTF_8)); }

游标(cursor)在创建自己的持久化 ledger 时调用该方法,将游标名写入pulsar/cursor(见 ManagedCursorImpl.java)。这样,游标 ledger 既带有 managed ledger 的基础属性,又额外标明自己属于哪个游标。

3. 压缩 ledger 元数据

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, // component=compacted-ledger METADATA_PROPERTY_COMPACTEDTOPIC, compactedTopic.getBytes(StandardCharsets.UTF_8), METADATA_PROPERTY_COMPACTEDTO, compactedToMessageId ); }

topic 压缩分两阶段执行,第二阶段(phaseTwo)在创建压缩 ledger 时调用此方法,传入原始 topic 名与to(压缩后的最后消息 id)序列化后的字节数组(见 TwoPhaseCompactor.java)。这正是原文档表格中pulsar/compactedTopic与pulsar/compactedTo两行数据的真实来源。

4. Schema ledger 元数据

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, // component=schema METADATA_PROPERTY_SCHEMAID, schemaId.getBytes(StandardCharsets.UTF_8) ); }

Broker 侧的 Schema 存储(BookkeeperSchemaStorage)在createLedger(String schemaId)中调用该方法构造元数据,并随bookKeeper.asyncCreateLedger(...)的metadata参数一并写入 ZooKeeper(见 BookkeeperSchemaStorage.java)。

5. 放置策略配置元数据

此外,当前版本还支持通过buildMetadataForPlacementPolicyConfig将EnsemblePlacementPolicyConfig编码进 ledger 元数据(键名为EnsemblePlacementPolicyConfig,见 LedgerMetadataUtils.java 与 EnsemblePlacementPolicyConfig.java),用于在 topic 级别定制 ledger 的放置策略。

五、如何实际读取元数据

由于元数据随 ledger 一起持久化在 ZooKeeper 中,读取方式与普通 BookKeeper 客户端一致:打开目标 ledger 后,通过LedgerHandle.getLedgerMetadata()拿到LedgerMetadata对象,再读取其自定义属性。Pulsar 内部也复用了这一机制——例如ManagedLedgerImpl.getLedgerMetadata(ledgerId)返回rh.getLedgerMetadata().toSafeString()的 JSON 文本(见 ManagedLedgerImpl.java),其中便包含上文所有的application、component、pulsar/...键值对。

以压缩场景为例,CompactedTopicTest中的测试直接展示了打开压缩 ledger 并校验其元数据属性的流程:先用bk.createLedger(...)创建 ledger,再通过bk.openLedger(ledgerId, ...)打开并读取(见 CompactedTopicTest.java)。实际运维排查时,可参照同样的思路:

  1. 通过pulsar-admin topics stats-internal拿到 topic 内部统计中的 ledger id 列表(getManagedLedgerInternalStats支持includeLedgerMetadata=true直接返回元数据文本,见 ManagedLedgerImpl.java);
  2. 用 BookKeeper 客户端(或bookkeeper shell ledgermetadata <ledgerId>)打开对应 ledger;
  3. 读取LedgerMetadata.getCustomMetadata()/toSafeString()输出,按本文表格中的键名对照解读。

六、元数据的运维实践价值

理解这些元数据后,你可以获得以下实际的排查与运维能力:

  • 快速识别 ledger 归属:component字段直接告诉你一个陌生 ledger 是 managed-ledger、schema 还是压缩产物,无需猜测;
  • 定位主题数据:pulsar/managed-ledger把 ledger 与具体 topic(ledger 名)精确对应,方便做存储分布统计与数据迁移核对;
  • 追踪游标状态:pulsar/cursor让游标持久化 ledger 与其消费游标一一对应,可用于排查游标堆积、回溯消费位置;
  • 审计压缩结果:pulsar/compactedTopic与pulsar/compactedTo记录了压缩覆盖的原始主题与最后压缩位置,可验证压缩任务是否按预期完成;
  • 识别 Schema 存储:较新版本中pulsar/schemaId帮助区分 Schema 专用 ledger。

需要注意的是:ledger 元数据在创建时一次性写入 ZooKeeper,属于静态描述信息;它与消息条数、字节大小等运行时统计不同,适合作为"这是什么"的定性依据,而不适合作为实时监控指标。结合getManagedLedgerInternalStats输出的 ensemble、quorum 等存储布局信息一起分析,可以获得对 Pulsar 存储层更完整的认知。

七、小结

本文完整覆盖了官方 Cookbook《BookKeeper Ledger Metadata》的全部内容:元数据存放于 ZooKeeper、可通过 BookKeeper API 读取,并详细列出application、component、pulsar/managed-ledger、pulsar/cursor、pulsar/compactedTopic、pulsar/compactedTo六类键的语义。在此基础上,我们从当前仓库源码 LedgerMetadataUtils.java 出发,还原了每一类元数据的构造时机与调用链(ManagedLedger 创建、游标持久化、两阶段压缩、Schema 存储),并给出了可落地的读取与排查方法。掌握这套元数据体系,你就能在 Pulsar 存储层排查中多一把精准的"放大镜"。

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

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

项目地址:https://gitcode.com/gh_mirrors/pulsar28/pulsar
点击查看免费下载

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

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询