- 消息队列
- 流处理
- 后端
- 微服务
- 消息路由
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
导读
本文围绕 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 的路径:
- 客户端(Client)路径:
org.apache.pulsar.client.admin.Topics接口提供基于MessageId的 Offload API,即 triggerOffload(String topic, MessageId messageId),用户必须显式指定一个 MessageId 作为卸载分界点,将其之前的所有数据卸载到冷存储。 - 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_VALUE | null | 数据总量远小于阈值,无可卸载数据 |
0 | MessageIdImpl(2, 0, -1) | 阈值极小,几乎全部数据都要卸载,卸载点取最新 ledger 的前一个 |
1000 | MessageIdImpl(2, 0, -1) | 累计 3000 > 1000,边界落在 ledger 2,卸载点为其前一个 ledger 1(id=2 是 reversed 顺序中的前一个,即原始最新的 ledger) |
5000 | MessageIdImpl(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); } }这里有三个实战要点值得注意:
- 大小阈值支持人类可读单位:
-s/--size-threshold参数通过ByteUnitToLongConverter转换,可直接书写10M、5G等形式,客户端 API 中则以纯字节数(long)传入。 - 当前 ledger 大小需手工补齐:
getInternalStats返回的currentLedgerSize不会自动填充到最后一个 ledger 的size字段(源码注释 "doesn't get filled in now it seems"),CLI 通过ledgers.get(ledgers.size() - 1).size = stats.currentLedgerSize显式补齐,否则最后一段数据量会被漏算。 - 空 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 调用链,理解这条链路有助于评估新方法的实际行为与限制:
- 客户端发起请求:
TopicsImpl.triggerOffloadAsync向admin/v2的 Topic 路径{topic}/offload发送PUT请求,请求体为MessageIdImpl(见 TopicsImpl.java)。 - Broker 校验与路由:
PersistentTopicsBase.internalTriggerOffload依次执行validateTopicOperationAsync(topicName, TopicOperation.OFFLOAD)(鉴权)、validateTopicOwnershipAsync(Topic 归属校验,非 owner 节点会收到 307 重定向)和getTopicReferenceAsync,最终调用((PersistentTopic) topic).triggerOffload(messageId)(见 PersistentTopicsBase.java)。 - PersistentTopic 执行:
PersistentTopic.triggerOffload是synchronized方法,若上一次卸载尚未完成则抛AlreadyRunningException(Broker 层转为 409 CONFLICT);否则基于换算出的MessageIdImpl构造Position,调用getManagedLedger().asyncOffloadPrefix(...)真正把该位置之前的数据写入长期存储(见 PersistentTopic.java)。 - 状态查询:卸载是异步长任务,可通过既有的
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 历史数据的能力从命令行下沉为正式客户端接口。其核心贡献在于:
- 明确了一个可复用的换算算法
findFirstLedgerWithinThreshold(已有 CLI 实现与单元测试背书); - 定义了
Topics.triggerOffload(String, long)与triggerOffloadAsync(String, long)两个新方法签名; - 通过零 Broker 改动实现完全向后兼容,降低了落地成本。
对于关注 Pulsar 存储成本治理的开发者,该接口使"为每个 Topic 设置 BookKeeper 保留上限、超额数据自动下沉冷存储"的运维策略可以直接在 Java 应用中编程化实现。
- 消息队列
- 流处理
- 后端
- 微服务
- 消息路由
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
PIP-348 源码级解析:Apache Pulsar 在 Topic 加载阶段触发分层存储 Offload
PIP 348 源码级解析:Apache Pulsar 在 Topic 加载阶段触发分层存储 Offload 导读 PIP 348(Trigger offloa
消息队列后端go-metrics 实践指南:为 distribution 构建规范化、可检索的 Prometheus 指标体系
go metrics 实践指南:为 distribution 构建规范化、可检索的 Prometheus 指标体系 go metrics 是 Docker 系列
消息队列流处理后端微服务消息路由Apache Pulsar PIP-348:在 Topic 加载阶段触发 Offload,让分层存储冷数据搬迁不再“等下一次”
Apache Pulsar PIP 348:在 Topic 加载阶段触发 Offload,让分层存储冷数据搬迁不再“等下一次” 本文基于 Apache Puls
消息队列流处理后端微服务消息路由
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考