- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
本指南以 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) |
|---|---|---|
| 所有 ledger | application | 'pulsar' |
| 所有 ledger | component | 'managed-ledger'、'schema'、'compacted-topic' |
| Managed ledgers | pulsar/managed-ledger | ledger 的名称(name of the ledger) |
| Cursor | pulsar/cursor | 游标名称(name of the cursor) |
| Compacted topic | pulsar/compactedTopic | 原始主题名称(name of the original topic) |
| Compacted topic | pulsar/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_APPLICATION | application | 应用标识,固定为pulsar |
METADATA_PROPERTY_COMPONENT | component | 组件标识 |
METADATA_PROPERTY_MANAGED_LEDGER_NAME | pulsar/managed-ledger | managed ledger 名称 |
METADATA_PROPERTY_CURSOR_NAME | pulsar/cursor | 游标名称 |
METADATA_PROPERTY_COMPACTEDTOPIC | pulsar/compactedTopic | 被压缩的原始主题 |
METADATA_PROPERTY_COMPACTEDTO | pulsar/compactedTo | 最后一条已压缩消息 id |
METADATA_PROPERTY_SCHEMAID | pulsar/schemaId | Schema 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)。实际运维排查时,可参照同样的思路:
- 通过
pulsar-admin topics stats-internal拿到 topic 内部统计中的 ledger id 列表(getManagedLedgerInternalStats支持includeLedgerMetadata=true直接返回元数据文本,见 ManagedLedgerImpl.java); - 用 BookKeeper 客户端(或
bookkeeper shell ledgermetadata <ledgerId>)打开对应 ledger; - 读取
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
相关推荐
Apache Pulsar BookKeeper Ledger 元数据完全指南:如何通过 ZooKeeper 与 BookKeeper API 解读数据存储结构
Apache Pulsar BookKeeper Ledger 元数据完全指南:如何通过 ZooKeeper 与 BookKeeper API 解读数据存储结构
消息队列后端流处理Apache Pulsar BookKeeper Ledger 元数据全解析:如何从 ZooKeeper 中读懂数据存储结构
Apache Pulsar BookKeeper Ledger 元数据全解析:如何从 ZooKeeper 中读懂数据存储结构 Apache Pulsar 的所有
消息队列后端流处理Apache Pulsar 配置完全指南:从 BookKeeper 到 ZooKeeper 的 conf 参数深度解析
Apache Pulsar 配置完全指南:从 BookKeeper 到 ZooKeeper 的 conf 参数深度解析 Apache Pulsar 是一个分布式
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考