news 2026/10/9 2:15:58

Apache Pulsar PIP-416 深度解读:基于存储大小阈值触发 Topic 数据 Offload 的客户端新接口

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Apache Pulsar PIP-416 深度解读:基于存储大小阈值触发 Topic 数据 Offload 的客户端新接口
  • 消息队列
  • 流处理
  • 后端
  • 微服务
  • 消息路由

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

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

导读

本文围绕 Apache Pulsar 社区的 PIP-416(Proposal for Improvement)展开,解读一项针对客户端管理接口的增强:在不改动 Broker REST API 的前提下,为Topics管理接口新增基于存储大小阈值触发 Topic 数据 Offload(卸载至长期存储)的方法。读完本文,你将掌握 PIP-416 的设计动机、大小阈值到 MessageId 的换算算法、同步与异步 API 的完整签名,以及从 CLI 命令到 Broker 底层 ManagedLedger 的完整调用链,可直接用于编写基于大小阈值的冷热数据分层管理程序。

背景:Pulsar 的分层存储与 Offload 机制

Apache Pulsar 支持将 BookKeeper 中的历史数据卸载(Offload)到长期存储(Long-term Storage),例如 AWS S3、GCS 等对象存储,从而降低热存储成本。在 PIP-416 之前,仓库中已经存在两条触发 Offload 的路径:

  1. 客户端(Client)路径:org.apache.pulsar.client.admin.Topics接口提供基于MessageId的 Offload API,即 triggerOffload(String topic, MessageId messageId),用户必须显式指定一个 MessageId 作为卸载分界点,将其之前的所有数据卸载到冷存储。
  2. CLI 路径:pulsar-admin命令行的topics offload命令支持基于存储大小阈值触发 Offload,用户只需给出"Topic 在 BookKeeper 中最多保留多少数据"(如10M、5G),CLI 内部会自动将其换算为一个 MessageId 再调用 Broker 接口。

PIP-416 的核心诉求,就是把 CLI 已经具备的"按大小阈值触发"能力下沉为正式的客户端 API,让 Java 程序可以直接以存储大小为粒度管理 Topic 的数据分层,而不必先查询内部统计、手工推算 MessageId。

动机:用户关注的是存储大小,而不是 MessageId

现有客户端 Offload 方法要求用户指定特定 MessageId 作为卸载点。但在真实生产场景中,用户通常更关心存储大小而非具体消息 ID——例如"这个 Topic 在 BookKeeper 中只保留最近 100GB 数据,更早的全部挪到 S3"。用户希望基于大小阈值触发 Offload,自动将一定量的历史数据移动到冷存储,而无需理解 ledger 与 MessageId 的内部细节。

PIP-416 的目标由此明确:

提供一个新的基于大小阈值的客户端 Topic 卸载方法,使用户能够更方便地管理 Topic 存储。

其 Scope 界定为:允许客户端通过指定存储大小阈值来触发 Topic 数据卸载,其余行为与既有 Offload 语义保持一致。

关键设计决策:复用 CLI 算法,不新增 Broker REST API

PIP-416 明确指出:无需为 Broker 的 REST API 新增接口,实现将参考 CLI(org.apache.pulsar.admin.cli)中Offload命令的做法,即先把sizeThreshold换算为具体 MessageId,再调用PersistentTopics中既有的triggerOffloadAPI。这样既复用了经过验证的换算逻辑,又保持了 Broker 端面接口的稳定。

在 CmdTopics.java 中可以找到该算法的完整实现(与 PIP-416 文档一致):

static MessageId findFirstLedgerWithinThreshold(List<PersistentTopicInternalStats.LedgerInfo> ledgers, long sizeThreshold) { long suffixSize = 0L; ledgers = Lists.reverse(ledgers); long previousLedger = ledgers.get(0).ledgerId; for (PersistentTopicInternalStats.LedgerInfo l : ledgers) { suffixSize += l.size; if (suffixSize > sizeThreshold) { return new MessageIdImpl(previousLedger, 0L, -1); } previousLedger = l.ledgerId; } return null; }

算法原理解读

  • 输入:Topic 的内部统计PersistentTopicInternalStats.LedgerInfo列表(每个 ledger 含ledgerId、entries、size),以及用户指定的大小阈值sizeThreshold(字节)。
  • 方向:Lists.reverse将 ledger 列表倒序,即从最新的 ledger 开始向最旧方向遍历,便于计算"从尾部累计的最近数据量"。
  • 核心逻辑:维护一个后缀和suffixSize,从最新的 ledger 开始累加其size;当累计大小首次超过阈值时,说明"从最新数据往前保留阈值大小的数据"这个边界落在当前 ledger 内,此时返回previousLedger(即当前 ledger 的前一个、更旧的 ledger)起始位置new MessageIdImpl(previousLedger, 0L, -1)作为卸载点——即该更旧 ledger 第 0 条 entry 之前的全部数据都应被卸载。
  • 返回 null:若遍历完所有 ledger,后缀和仍不超过阈值(即整个 Topic 数据量都不足阈值),说明"没有需要卸载的数据",返回null。

在 TestCmdTopics.java 中,有针对该算法的单元测试testFindFirstLedgerWithinThreshold,构造了三个 ledger(ledger 0: 1000 bytes、ledger 1: 2000 bytes、ledger 2: 3000 bytes),验证:

阈值期望结果说明
Long.MAX_VALUEnull数据总量远小于阈值,无可卸载数据
0MessageIdImpl(2, 0, -1)阈值极小,几乎全部数据都要卸载,卸载点取最新 ledger 的前一个
1000MessageIdImpl(2, 0, -1)累计 3000 > 1000,边界落在 ledger 2,卸载点为其前一个 ledger 1(id=2 是 reversed 顺序中的前一个,即原始最新的 ledger)
5000MessageIdImpl(1, 0, -1)累计到 ledger 1 时 2000 + 3000 > 5000,卸载点为 ledger 1 的前一个 ledger 0

该测试同时佐证了 CLI 与 PIP-416 客户端方法将共享的换算语义。

CLI Offload 命令的完整上下文

为了让换算结果可用,CLI 的Offload命令在调用换算前还会做一步关键补齐,见 CmdTopics.java:

@Command(description = "Trigger offload of data from a topic to long-term storage (e.g. Amazon S3)") private class Offload extends CliCommand { @Option(names = { "-s", "--size-threshold" }, description = "Maximum amount of data to keep in BookKeeper for the specified topic (e.g. 10M, 5G).", required = true, converter = ByteUnitToLongConverter.class) private Long sizeThreshold; @Parameters(description = "persistent://tenant/namespace/topic", arity = "1") private String topicName; @Override void run() throws PulsarAdminException { String persistentTopic = validatePersistentTopic(topicName); PersistentTopicInternalStats stats = getTopics().getInternalStats(persistentTopic, false); if (stats.ledgers.size() < 1) { throw new PulsarAdminException("Topic doesn't have any data"); } LinkedList<PersistentTopicInternalStats.LedgerInfo> ledgers = new LinkedList<>(stats.ledgers); ledgers.get(ledgers.size() - 1).size = stats.currentLedgerSize; // doesn't get filled in now it seems MessageId messageId = findFirstLedgerWithinThreshold(ledgers, sizeThreshold); if (messageId == null) { System.out.println("Nothing to offload"); return; } getTopics().triggerOffload(persistentTopic, messageId); System.out.println("Offload triggered for " + persistentTopic + " for messages before " + messageId); } }

这里有三个实战要点值得注意:

  1. 大小阈值支持人类可读单位:-s/--size-threshold参数通过ByteUnitToLongConverter转换,可直接书写10M、5G等形式,客户端 API 中则以纯字节数(long)传入。
  2. 当前 ledger 大小需手工补齐:getInternalStats返回的currentLedgerSize不会自动填充到最后一个 ledger 的size字段(源码注释 "doesn't get filled in now it seems"),CLI 通过ledgers.get(ledgers.size() - 1).size = stats.currentLedgerSize显式补齐,否则最后一段数据量会被漏算。
  3. 空 Topic 直接报错:stats.ledgers.size() < 1时抛出PulsarAdminException("Topic doesn't have any data");换算结果为null时打印 "Nothing to offload" 并直接返回,不会发起无意义的 Offload 请求。

PIP-416 的客户端实现需要完整继承这套语义:先取内部统计、补齐当前 ledger 大小、换算 MessageId、判断空结果,再落到既有triggerOffload调用。

新增公共 API:Topics 接口

PIP-416 在org.apache.pulsar.client.admin.Topics接口中新增两个方法声明(同步与异步各一),完整代码来自 PIP-416 文档:

/** * Trigger offload of data to long-term storage based on size threshold * * @param topic * Topic name * @param sizeThreshold * Size threshold in bytes * @throws PulsarAdminException */ void triggerOffload(String topic, long sizeThreshold) throws PulsarAdminException; /** * Trigger offload of data to long-term storage based on size threshold asynchronously * * @param topic * Topic name * @param sizeThreshold * Size threshold in bytes * @return Future that completes once the offload operation has started */ CompletableFuture<Void> triggerOffloadAsync(String topic, long sizeThreshold);

与既有 triggerOffload(String, MessageId) 相比,新方法的差异点在于:

  • 第二参数从MessageId变为long sizeThreshold(字节),调用方无需感知 ledger/entry 结构;
  • 语义为"触发卸载",返回的CompletableFuture<Void>在 Offload 操作开始时即完成(而非等待卸载结束),与既有 API 的异步契约保持一致;
  • 同步方法包装异步方法(类似 TopicsImpl.triggerOffload 中sync(() -> triggerOffloadAsync(topic, messageId))的既有模式),抛出的PulsarAdminException由同步包装层统一转换。

Broker 端底层调用链:从换算到实际卸载

虽然 PIP-416 不新增 Broker REST API,但客户端新方法最终仍会落入既有的 Broker 调用链,理解这条链路有助于评估新方法的实际行为与限制:

  1. 客户端发起请求:TopicsImpl.triggerOffloadAsync向admin/v2的 Topic 路径{topic}/offload发送PUT请求,请求体为MessageIdImpl(见 TopicsImpl.java)。
  2. Broker 校验与路由:PersistentTopicsBase.internalTriggerOffload依次执行validateTopicOperationAsync(topicName, TopicOperation.OFFLOAD)(鉴权)、validateTopicOwnershipAsync(Topic 归属校验,非 owner 节点会收到 307 重定向)和getTopicReferenceAsync,最终调用((PersistentTopic) topic).triggerOffload(messageId)(见 PersistentTopicsBase.java)。
  3. PersistentTopic 执行:PersistentTopic.triggerOffload是synchronized方法,若上一次卸载尚未完成则抛AlreadyRunningException(Broker 层转为 409 CONFLICT);否则基于换算出的MessageIdImpl构造Position,调用getManagedLedger().asyncOffloadPrefix(...)真正把该位置之前的数据写入长期存储(见 PersistentTopic.java)。
  4. 状态查询:卸载是异步长任务,可通过既有的offloadStatus/offloadStatusAsync查询OffloadProcessStatus(NOT_RUN / RUNNING / SUCCESS / ERROR)。

从源码结构看,PIP-416 的新方法将复用这条链路,差异仅在于客户端本地完成"大小阈值 → MessageId"的换算,因此对 Broker 而言与既有 MessageId 触发的卸载完全等价,这也是 PIP-416 声称"完全向后兼容"的底气所在。

兼容性分析

PIP-416 明确说明:完全兼容(Fully compatible)。这是一项新的客户端方法实现,不影响既有功能,具体体现在:

  • 不修改、不删除任何既有 API 签名;
  • 不新增 Broker REST 接口,Broker 端无需升级即可配合新客户端工作(前提是换算后的 MessageId 语义与 Broker 期望一致);
  • 新方法与既有 MessageId 版本并存,两者都通过同一{topic}/offload端点触发卸载;
  • 卸载目标长期存储的驱动(S3、GCS 等)与分层存储策略(见 conf/broker.conf 与 conf/standalone.conf 中的 offload 相关配置)均由既有机制决定,不受新接口影响。

实战:基于大小阈值的客户端卸载示例

基于 PIP-416 的接口设计,Java 客户端使用方式形如(示意代码,需结合PulsarAdmin实例与异常处理):

PulsarAdmin admin = PulsarAdmin.builder().serviceHttpUrl("http://broker.example:8080").build(); String topic = "persistent://public/default/orders"; // 同步触发:将 Topic 在 BookKeeper 中保留的数据压缩到 100GB 以内,更早的数据卸载到长期存储 try { admin.topics().triggerOffload(topic, 100L * 1024 * 1024 * 1024); System.out.println("Offload triggered (threshold=100GB)"); } catch (PulsarAdminException e) { // 处理鉴权失败、Topic 无数据、卸载已在运行(409)等异常 } // 异步触发:不阻塞调用线程 admin.topics().triggerOffloadAsync(topic, 100L * 1024 * 1024 * 1024) .thenRun(() -> System.out.println("Offload operation has started")); // 查询卸载进度 OffloadProcessStatus status = admin.topics().offloadStatus(topic);

实战注意事项:

  • 阈值单位是字节:客户端 API 接收long字节数,与 CLI 的10M/5G可读写法不同,需要自行换算;
  • 阈值语义是"BookKeeper 中保留的最大数据量":即从最新数据往回数,累计超过阈值的历史数据会被卸载,与 CLI 中--size-threshold的语义一致;
  • 异步 Future 完成仅代表"卸载已开始":真正的卸载完成需轮询offloadStatus,可以参考 CLI 中 OffloadStatusCmd 的--wait-complete轮询模式(每秒查询一次直至非 RUNNING 状态);
  • Topic 必须为持久化 Topic:Offload 仅适用于persistent://域下的 Topic,validatePersistentTopic会校验这一点。

总结

PIP-416 以"复用 CLI 换算算法、复用 Broker 既有触发链路、新增客户端 API"三个动作,把按存储大小阈值卸载 Topic 历史数据的能力从命令行下沉为正式客户端接口。其核心贡献在于:

  1. 明确了一个可复用的换算算法findFirstLedgerWithinThreshold(已有 CLI 实现与单元测试背书);
  2. 定义了Topics.triggerOffload(String, long)与triggerOffloadAsync(String, long)两个新方法签名;
  3. 通过零 Broker 改动实现完全向后兼容,降低了落地成本。

对于关注 Pulsar 存储成本治理的开发者,该接口使"为每个 Topic 设置 BookKeeper 保留上限、超额数据自动下沉冷存储"的运维策略可以直接在 Java 应用中编程化实现。

  • 消息队列
  • 流处理
  • 后端
  • 微服务
  • 消息路由

【免费下载链接】pulsar

Apache Pulsar - distributed pub-sub messaging system

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

相关推荐

上一篇:百度网盘秒传链接网页工具终极指南:从零开始快速掌握文件极速转存
下一篇:百度网盘秒传链接网页工具完全指南:跨平台免费极速转存教程

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

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

Claude Code九成不好用?多半是用错了方式

十个说Claude Code不好用的人里&#xff0c;有九个是把它用错了地方。我第一次接触Claude Code时&#xff0c;心里想的就是"这不就是个跑在终端里的ChatGPT吗"——问它怎么改代码&#xff0c;让它写个函数&#xff0c;再把结果复制回编辑器&#xff0c;试用两天后我得…

作者头像 李华
网站建设 2026/10/9 2:12:35

RAG 入门与实践指南:用 TaoToken 统一 Key 打通检索增强生成全链路

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

作者头像 李华
网站建设 2026/10/9 2:12:29

Python眼底图像视盘视杯分割实战:从U-Net到CDR指标计算

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

作者头像 李华