消息队列流处理后端微服务消息路由【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pu/pulsar点击查看免费下载导读Apache Pulsar 传统的 backlog 统计只给出 entries存储条目数量而消费者实际感受到的工作量与消息messages数量成正比——尤其在批处理batching与服务端过滤器PIP-105 Subscription Filters共同作用的场景下entries 与 messages 严重不成比例导致积压监控失真。PIP-187 为 Pulsar 引入了一套全新的 订阅 backlog 精确分析 能力通过新增 ManagedCursor 扫描 API、Broker REST 接口、Java PulsarAdmin 方法与 pulsar-admin CLI 命令从指定位置默认最近未确认位置 lastMarkDeletePosition开始逐条读取存储中的消息并施加订阅过滤器最终返回 entries、messages、过滤器裁决结果等精确统计。本文将以 pip/pip-187.md 为骨架结合仓库源码与测试用例完整讲解该特性的 API 形态、配置参数、底层实现原理、客户端循环扫描机制与最佳实践。问题背景为什么传统 backlog 统计不够精确PIP-187 提案明确指出在 PIP-187 落地之前Pulsar 无法为一个订阅提供精确的 backlog 数值原因集中在两点只有 entries 数量没有 messages 数量。Pulsar 底层存储ManagedLedger以 entry 为单位组织数据而单个 entry 可以承载一条或多条消息。当生产者开启批处理batching时一个 entry 内可能打包多条消息消费者侧的实际处理工作量与消息条数成正比与 entry 条数不成正比——仅靠 entry 数无法估算真实的消费压力。服务端过滤器PIP-105可能过滤掉部分消息。PIP-105 引入了订阅级别的服务端过滤器Subscription Filters经过过滤器 REJECT/RESCHEDULE 的消息实际上不会投递给消费者但它们依然占用存储和 entry 计数导致传统统计口径下的 backlog 虚高。此外由于订阅具备极强的动态性订阅可以从 earliest 位置创建、过滤器可以随时间变化这些信息无法通过简单的计数器在写入路径上预先维护唯一可靠的精确值只能通过事后扫描计算得出。这正是 analyzeBacklog 系列 API 存在的根本原因。API 变更总览PIP-187 的实现横跨了存储层、Broker 管理接口与客户端工具链四个层次层次新增内容说明存储层ManagedCursorscan(OptionalPosition, PredicateEntry, ...)从指定位置扫描游标到尾部带时间与条目上限Broker RESTPOST .../subscription/{subName}/analyzeBacklog异步接口返回精确分析结果Java PulsarAdminanalyzeSubscriptionBacklog(...)系列重载同步/异步、支持客户端循环扫描pulsar-admin CLIpulsar-admin topics analyze-backlog命令行入口支持进度输出与 NDJSON下面逐层展开。存储层ManagedCursor.scan APIPIP-187 在内部 ManagedCursor 接口中新增了扫描能力原始提案给出的签名是CompletableFutureScanOutcome scan(OptionalPosition startingPosition, PredicateEntry condition, long maxEntries, long timeOutMs);在仓库中的最终实现位于 managed-ledger/src/main/java/org/apache/bookkeeper/mledger/ManagedCursor.java签名在此基础上增加了batchSize参数用于控制每次从存储读取的条目数default CompletableFutureScanOutcome scan(OptionalPosition startingPosition, PredicateEntry condition, int batchSize, long maxEntries, long timeOutMs) { return CompletableFuture.failedFuture(new UnsupportedOperationException()); }关键语义如下起点startingPosition未提供时从lastMarkDeletePosition最近一次 mark-delete 位置即已确认消息的边界开始提供时从该位置开始。终止条件扫描受maxEntries最大条目数与timeOutMs最大耗时双重限制目的是防止在积压巨大时产生巨大且无用的全量扫描起到保护 Broker 的作用。条件回调PredicateEntry逐条检查 entry返回false可提前终止扫描返回的ScanOutcome枚举则向调用方传达本次扫描的最终状态COMPLETED表示正常扫完否则视为被中断/中止。成本警告接口 Javadoc 明确标注 this is an expensive operation因为它需要从存储BookKeeper读取每条消息并执行过滤逻辑。Broker REST 接口端点与请求形态提案中的 REST 原型为 GET QueryParam 传 position但在仓库中的最终实现pulsar-broker/src/main/java/org/apache/pulsar/broker/admin/v2/PersistentTopics.java演进为POST RequestBody 传 positionPOST /admin/v2/persistent/{tenant}/{namespace}/{topic}/subscription/{subName}/analyzeBacklog参数说明参数位置必填说明tenant/namespace/topicPath是主题的完整定位信息topic需 URL 编码subNamePath是被分析的订阅名需 URL 编码positionRequestBodyResetCursorData含ledgerId、entryId否扫描起点不传则从lastMarkDeletePosition开始authoritativeQuery否默认 false是否为 leader broker 重定向到本 broker 的调用内部使用服务端会先校验主题所有权非 authoritative 场景下若本 Broker 不服务该 namespace会返回 307 重定向到正确的 Broker再校验订阅是否存在不存在返回 404以及调用方是否具备该主题的CONSUME权限无权限返回 401/403。响应模型AnalyzeSubscriptionBacklogResult返回的 JSON 对象对应 pulsar-client-admin-api/src/main/java/org/apache/pulsar/common/stats/AnalyzeSubscriptionBacklogResult.java字段如下字段类型含义entrieslong扫描到的存储条目总数messageslong扫描到的逻辑消息总数含批处理展开markerMessageslong扫描到的内部 marker 消息数如事务标记不计入待投递消息filterAcceptedEntrieslong被过滤器 ACCEPT 的条目数filterRejectedEntrieslong被过滤器 REJECT 的条目数filterRescheduledEntrieslong被过滤器 RESCHEDULE 的条目数filterAcceptedMessageslong被过滤器 ACCEPT 的消息数filterRejectedMessageslong被过滤器 REJECT 的消息数filterRescheduledMessageslong被过滤器 RESCHEDULE 的消息数abortedboolean是否因内部限制超时、条目上限被中止firstMessageIdString扫描到的第一条消息 IDledgerId:entryIdlastMessageIdString扫描到的最后一条消息 IDledgerId:entryId其中abortedtrue表示请求被某些内部限制如超时或达到最大条目数中止。PIP-187 明确说明API 不会提供中止原因的更多细节避免接口过于繁琐、难以维护如需排查Broker 日志中会有详细记录。firstMessageId/lastMessageId是最终实现相对提案的增强它们为客户端循环扫描见下文提供了续扫锚点。Java PulsarAdmin API同步与异步入口接口定义位于 pulsar-client-admin-api/src/main/java/org/apache/pulsar/client/admin/Topics.java实现位于 pulsar-client-admin/src/main/java/org/apache/pulsar/client/admin/internal/TopicsImpl.java。核心方法族如下// 同步从指定起点empty 表示最后已处理位置开始分析 AnalyzeSubscriptionBacklogResult analyzeSubscriptionBacklog(String topic, String subscriptionName, OptionalMessageId startPosition) throws PulsarAdminException; // 同步额外指定客户端侧最大扫描条目数用于客户端循环终止 AnalyzeSubscriptionBacklogResult analyzeSubscriptionBacklog(String topic, String subscriptionName, OptionalMessageId startPosition, long backlogScanMaxEntries) throws PulsarAdminException; // 同步自定义终止谓词 AnalyzeSubscriptionBacklogResult analyzeSubscriptionBacklog(String topic, String subscriptionName, OptionalMessageId startPosition, PredicateAnalyzeSubscriptionBacklogResult terminatePredicate) throws PulsarAdminException; // 异步版本三个重载一一对应 CompletableFutureAnalyzeSubscriptionBacklogResult analyzeSubscriptionBacklogAsync(String topic, String subscriptionName, OptionalMessageId startPosition); CompletableFutureAnalyzeSubscriptionBacklogResult analyzeSubscriptionBacklogAsync(String topic, String subscriptionName, OptionalMessageId startPosition, long backlogScanMaxEntries); CompletableFutureAnalyzeSubscriptionBacklogResult analyzeSubscriptionBacklogAsync(String topic, String subscriptionName, OptionalMessageId startPosition, PredicateAnalyzeSubscriptionBacklogResult terminatePredicate);topic必须是持久主题persistent://...startPosition为Optional.empty()时表示从订阅最后已处理位置mark-delete 位置开始扫描。所有方法均为潜在昂贵操作Javadoc 反复强调其会从存储读取消息、并考虑批处理消息与订阅过滤器。客户端循环扫描Client-side Loop机制Broker 侧有subscriptionBacklogScanMaxEntries默认 10000与subscriptionBacklogScanMaxTimeMs默认 2 分钟两个硬上限积压超过上限时单次请求必然abortedtrue。为解决想拿到完整精确积压与Broker 单次扫描有上限的矛盾客户端实现了自动循环续扫客户端先发起一次analyze-backlog请求若服务端返回abortedtrue且未满足客户端终止条件客户端取出响应中的lastMessageId将其entryId 1作为新的起点TopicsImpl.java再次发起请求循环中每次返回的统计通过mergeBacklogResults累加合并直到服务端返回abortedfalse扫描完整或累计entries达到客户端设定的backlogScanMaxEntries或满足自定义terminatePredicate防御性兜底若某次返回entries 0或lastMessageId为空视为扫描完成终止循环TopicsImpl.java。需要特别强调两点语义源码 Javadoc 原话若服务端subscriptionBacklogScanMaxEntries 客户端 backlogScanMaxEntries则客户端参数不生效backlogScanMaxEntries并不能精确控制服务端扫描的条目数它只决定客户端循环何时停止——实际扫描的条目数会是服务端单次上限的整数倍。这一机制让运维人员可以在不调大 Broker 全局限制避免单次请求长期占用线程与 IO的前提下通过客户端循环拿到任意深的精确积压。pulsar-admin CLI 命令对应命令定义在 pulsar-client-tools/src/main/java/org/apache/pulsar/admin/cli/CmdTopics.javapulsar-admin topics analyze-backlog persistent://tenant/namespace/topic \ -s subscription-name \ [-p ledgerId:entryId] \ [-b backlogScanMaxEntries] \ [-q] \ [--plain]选项必填默认说明-s, --subscription是—待分析的订阅名-p, --position否最近未确认位置扫描起点格式ledgerId:entryId如-p 5:100也支持-1:-1从 earliest 起扫-b, --backlog-scan-max-entries否不限制客户端循环终止阈值必须大于 0否则报参数错误-q, --quiet否false关闭循环扫描时的进度输出--plain否false以 NDJSON每行一条 JSON而非美化 JSON 输出结果不指定-b时命令调用单次analyzeSubscriptionBacklog(topic, sub, startPosition)指定-b时以result.getEntries() backlogScanMaxEntries为终止谓词调用循环版本并且每次循环都会打印一次中间结果除非--quiet。最终结果统一通过print输出--plain便于脚本用 jq 等工具解析。Broker 配置项两个新增配置定义在 pulsar-broker-common/src/main/java/org/apache/pulsar/broker/ServiceConfiguration.java并同步出现在 conf/broker.conf# Maximum time in ms for a Analise backlog operation to complete subscriptionBacklogScanMaxTimeMs120000 # Maximum number of entries to be read within a Analise backlog operation subscriptionBacklogScanMaxEntries10000配置项默认值说明subscriptionBacklogScanMaxTimeMs1200002 分钟单次服务端扫描允许的最大耗时超时则返回abortedtruesubscriptionBacklogScanMaxEntries10000单次服务端扫描允许处理的最大条目数达到即返回abortedtrue注意调大这两个值可能让单次 HTTP 请求长时间挂起甚至先于命令完成前触发 HTTP 请求超时以及 NAT/防火墙的 idle timeout因此源码 Javadoc 给出的建议是尽量保持 Broker 侧限制收敛需要深扫时依靠客户端循环backlogScanMaxEntries来获取更多条目。此外扫描时每次读取的批大小复用dispatcherMaxReadBatchSize配置见下文实现细节。实现原理扫描过程在 Broker 内部如何运作核心执行路径PersistentSubscription.analyzeBacklog是整条链路的枢纽位于 pulsar-broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentSubscription.java调用链如下PersistentTopics.analyzeSubscriptionBacklog (REST) → PersistentTopicsBase.internalAnalyzeSubscriptionBacklogForNonPartitionedTopic → Subscription.analyzeBacklog(OptionalPosition) → PersistentSubscription.analyzeBacklog → cursor.duplicateNonDurableCursor(...) // 创建临时非持久游标 → newNonDurableCursor.scan(position, condition, batchSize, maxEntries, timeOutMs) → managedLedger.asyncDeleteCursor(...) // 扫描结束清理临时游标实现要点临时非持久游标扫描基于cursor.duplicateNonDurableCursor(analyze-backlog- UUID)创建的副本游标避免干扰订阅的真实消费游标扫描完成后异步删除该临时游标asyncDeleteCursor不留垃圾。从 mark-delete 位置起扫游标副本的getMarkDeletedPosition()即最后已处理位置未指定起点时从该处开始。逐条读取并解析元数据对每个 entry若内存中无元数据则通过Commands.peekMessageMetadata解析有markerType的消息内部 marker如事务标记单独计入markerMessages。批处理展开若messageMetadata.hasNumMessagesInBatch()则取getNumMessagesInBatch()作为该 entry 包含的逻辑消息数用于累计messages及各类 filter 消息计数。过滤器裁决通过entryFilterSupport.runFiltersForEntry(entry, messageMetadata, null)执行订阅过滤器EntryFilterSupport来自 dispatcher 或新建实例结果按EntryFilter.FilterResult分为ACCEPT/REJECT/RESCHEDULE三档分别计数对应响应模型中的三组 accepted/rejected/rescheduled 统计。硬性保护限制maxEntries与timeOutMs直接取自上述两个配置目的是防止拒绝服务——这与 PIP-187 中prevent huge (and useless) scans的初衷一致。进度日志每处理 1000 条 entry 打印一条Scan running日志含已耗时与条数结束后打印Scan complete含总耗时与结果供运维在 Broker 日志中定位中止原因。分区主题与不支持场景分区主题不可直接分析测试 AnalyzeBacklogSubscriptionTest.partitionedTopicNotAllowed 验证了对分区主题调用会抛出PulsarAdminException.NotAllowedException但可以分别对topic-partition-N单分区执行分析。非持久主题不支持NonPersistentSubscription.analyzeBacklog 直接抛出UnsupportedOperationException(Unsupported operation analyzeBacklog for NonPersistentSubscription)——该特性只针对持久订阅。与 individuallyDeletedMessages 的交互PIP-187 实现说明中专门提到它考虑 individuallyDeletedMessages单条确认的消息。scan基于游标副本扫描而游标的individuallyDeletedMessages单条删除范围会被扫描逻辑跳过因此已被单独确认的消息不会计入 backlog。测试 analyzeBacklogWithIndividualAck 与simpleAnalyzeBacklogTestL59-L151对这一行为做了验证例如非批处理场景下单条 ack 后entries与messages各减 1批处理场景下 ack 单条消息不减少 entry 计数entry 仍存在这与底层存储语义一致。测试验证从源码看行为保证仓库中的核心测试文件是 pulsar-broker/src/test/java/org/apache/pulsar/broker/admin/AnalyzeBacklogSubscriptionTest.java覆盖了以下关键场景测试方法验证点simpleAnalyzeBacklogTest/simpleAnalyzeBacklogTestWithBatching无批处理与批处理batchSize5下 entries/messages 的精确计数消费 ack 后统计随之减少partitionedTopicNotAllowed分区主题整体不可用单分区可用analyzeBacklogServerReturnFalseAbortedFlagWithoutLoop积压小于服务端上限时abortedfalseanalyzeBacklogMaxEntriesExceedWithoutLoop积压超过服务端上限时返回abortedtrue只统计到上限analyzeBacklogServerReturnFalseAbortedFlagWithLoop客户端循环下可拿到完整积压analyzeBacklogMaxEntriesExceedWithLoop客户端backlogScanMaxEntries控制循环终止Broker 单次 15 条、循环 3 次共返回 45 条客户端阈值 40analyzeBacklogWithTopicUnload扫描期间主题被 unload/重新加载后循环续扫仍正确analyzeBacklogWithIndividualAck单条 ackindividuallyDeletedMessages场景下计数正确此外tests/integration/src/test/java/org/apache/pulsar/tests/integration/cli/topic/AnalyzeBacklogTest.java 还提供了针对pulsar-admin topics analyze-backlog命令的集成测试。为何放弃旁路计数方案Rejected AlternativesPIP-187 在拒绝的备选方案一节中解释了为什么不能通过更廉价的方式获得精确值这对理解该 API 的定位至关重要写入时维护逻辑消息计数器。不可行的原因有四写入时无法预知未来会有哪些订阅订阅在写入之后才可能创建订阅可以从过去Earliest创建计数器无法回溯覆盖订阅过滤器通常配置在 Subscription Properties 上、是动态可变的计数器无法跟随过滤器变化在写入路径上执行过滤器计算会拖垮写入延迟与吞吐。用客户端克隆订阅并消费数据来计算。同样不可行这要求把海量积压数据从 Broker 传输到客户端工作量大且严重浪费带宽与资源。这两条理由也反向解释了为什么该功能必须是按需、扫描式、读路径的昂贵操作——这是权衡所有方案后唯一既能保证精确、又不对写入路径造成影响的可行路径。实践建议如何用它做监控与告警PIP-187 给出的监控策略非常明确可以作为生产环境落地该 API 的准则常规告警用 stats日常监控、自动告警应继续基于常规的订阅统计backlog entries 等配置阈值保持低开销精确分析走 analyzeBacklog当常规统计出现异常积压疑似超标、怀疑过滤器误过滤、需要评估消费延迟的真实工作量时通过pulsar-admin topics analyze-backlog或 Java API 发起一次精确分析作为人工介入的深度排查手段控制扫描成本由于该操作昂贵建议使用-b客户端循环限制或依赖 Broker 默认上限避免单次请求扫描过深对超大积压可先观察aborted标志与lastMessageId判断扫描是否完整配合过滤器排查对比filterRejectedMessages与filterAcceptedMessages可快速判断积压中多少消息实际不会投递从而校准真实消费压力注意适用边界仅持久主题、非分区主题或单分区可用Broker 配置subscriptionBacklogScanMaxTimeMs/subscriptionBacklogScanMaxEntries控制单次扫描的硬上限调整时需权衡 HTTP 请求超时风险。小结PIP-187 通过按需扫描 服务端过滤 批处理展开的组合为 Apache Pulsar 提供了唯一准确的订阅 backlog 视图ManagedCursor.scan提供存储层扫描原语Broker REST 与 Java PulsarAdmin 暴露分析入口客户端循环机制突破单次扫描上限CLI 命令让运维人员一条命令即可洞察积压真相。理解其昂贵但必要的设计取舍以及 entries/messages/过滤器裁决等多维统计口径是正确使用这一能力进行监控、告警与容量评估的关键。相关实现与测试均可直接在本仓库对应路径中查阅便于深入研读。赞分享消息队列流处理后端微服务消息路由【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pu/pulsar点击查看免费下载相关推荐Apache Pulsar PIP-33 复制订阅Replicated Subscriptions设计与实现解析Apache Pulsar PIP 33 复制订阅Replicated Subscriptions设计与实现解析 本篇围绕 pip/pip 33.md ht消息队列流处理后端微服务消息路由Apache Pulsar PIP-13 深度解析正则表达式主题订阅的实现原理与实战Apache Pulsar PIP 13 深度解析正则表达式主题订阅的实现原理与实战 本文围绕 Apache Pulsar 的 PIP 13Subscrib消息队列流处理后端微服务消息路由Apache Pulsar PIP-313 实战解析Consumer API 强制 Unsubscribe 的设计与实现Apache Pulsar PIP 313 实战解析Consumer API 强制 Unsubscribe 的设计与实现 本文围绕 PIP 313 https消息队列流处理后端微服务消息路由上一篇抖音无水印下载终极指南3分钟学会批量保存高清视频的免费工具下一篇抖音无水印下载终极指南3步轻松保存高清视频的免费工具创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考