
消息队列流处理后端微服务消息路由【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pu/pulsar点击查看免费下载导读本文围绕 Apache Pulsar 社区的 PIP-416Proposal 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、5GCLI 内部会自动将其换算为一个 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 APIPIP-416 明确指出无需为 Broker 的 REST API 新增接口实现将参考 CLIorg.apache.pulsar.admin.cli中Offload命令的做法即先把sizeThreshold换算为具体 MessageId再调用PersistentTopics中既有的triggerOffloadAPI。这样既复用了经过验证的换算逻辑又保持了 Broker 端面接口的稳定。在 CmdTopics.java 中可以找到该算法的完整实现与 PIP-416 文档一致static MessageId findFirstLedgerWithinThreshold(ListPersistentTopicInternalStats.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构造了三个 ledgerledger 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 1id2 是 reversed 顺序中的前一个即原始最新的 ledger5000MessageIdImpl(1, 0, -1)累计到 ledger 1 时 2000 3000 5000卸载点为 ledger 1 的前一个 ledger 0该测试同时佐证了 CLI 与 PIP-416 客户端方法将共享的换算语义。CLI Offload 命令的完整上下文为了让换算结果可用CLI 的Offload命令在调用换算前还会做一步关键补齐见 CmdTopics.javaCommand(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 doesnt have any data); } LinkedListPersistentTopicInternalStats.LedgerInfo ledgers new LinkedList(stats.ledgers); ledgers.get(ledgers.size() - 1).size stats.currentLedgerSize; // doesnt 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字段源码注释 doesnt get filled in now it seemsCLI 通过ledgers.get(ledgers.size() - 1).size stats.currentLedgerSize显式补齐否则最后一段数据量会被漏算。空 Topic 直接报错stats.ledgers.size() 1时抛出PulsarAdminException(Topic doesnt have any data)换算结果为null时打印 Nothing to offload 并直接返回不会发起无意义的 Offload 请求。PIP-416 的客户端实现需要完整继承这套语义先取内部统计、补齐当前 ledger 大小、换算 MessageId、判断空结果再落到既有triggerOffload调用。新增公共 APITopics 接口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 */ CompletableFutureVoid triggerOffloadAsync(String topic, long sizeThreshold);与既有 triggerOffload(String, MessageId) 相比新方法的差异点在于第二参数从MessageId变为long sizeThreshold字节调用方无需感知 ledger/entry 结构语义为触发卸载返回的CompletableFutureVoid在 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)鉴权、validateTopicOwnershipAsyncTopic 归属校验非 owner 节点会收到 307 重定向和getTopicReferenceAsync最终调用((PersistentTopic) topic).triggerOffload(messageId)见 PersistentTopicsBase.java。PersistentTopic 执行PersistentTopic.triggerOffload是synchronized方法若上一次卸载尚未完成则抛AlreadyRunningExceptionBroker 层转为 409 CONFLICT否则基于换算出的MessageIdImpl构造Position调用getManagedLedger().asyncOffloadPrefix(...)真正把该位置之前的数据写入长期存储见 PersistentTopic.java。状态查询卸载是异步长任务可通过既有的offloadStatus/offloadStatusAsync查询OffloadProcessStatusNOT_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 (threshold100GB)); } 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 必须为持久化 TopicOffload 仅适用于persistent://域下的 TopicvalidatePersistentTopic会校验这一点。总结PIP-416 以复用 CLI 换算算法、复用 Broker 既有触发链路、新增客户端 API三个动作把按存储大小阈值卸载 Topic 历史数据的能力从命令行下沉为正式客户端接口。其核心贡献在于明确了一个可复用的换算算法findFirstLedgerWithinThreshold已有 CLI 实现与单元测试背书定义了Topics.triggerOffload(String, long)与triggerOffloadAsync(String, long)两个新方法签名通过零 Broker 改动实现完全向后兼容降低了落地成本。对于关注 Pulsar 存储成本治理的开发者该接口使为每个 Topic 设置 BookKeeper 保留上限、超额数据自动下沉冷存储的运维策略可以直接在 Java 应用中编程化实现。赞分享消息队列流处理后端微服务消息路由【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pu/pulsar点击查看免费下载相关推荐PIP-348 源码级解析Apache Pulsar 在 Topic 加载阶段触发分层存储 OffloadPIP 348 源码级解析Apache Pulsar 在 Topic 加载阶段触发分层存储 Offload 导读 PIP 348Trigger 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),仅供参考