ARTICLE DETAIL

资讯详情

深耕网站建设与运营推广的一线实战洞察。

Apache Pulsar 主题统计(Topic Stats)详解:从命令行到源码级的指标解读

Apache Pulsar 主题统计(Topic Stats)详解:从命令行到源码级的指标解读 消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载Apache Pulsar 提供了丰富的主题级topic-level统计信息用于观测消息生产/消费速率、积压backlog、订阅状态以及跨机房复制replication情况。本文以 Pulsar 2.2.1 文档中的 administration-stats 为主体结合当前仓库的源码实现pulsar-common与pulsar-broker模块系统讲解分区主题partitioned topic与主题内部internal统计中每一个字段的真实含义、获取方式与底层原理帮助你快速定位生产环境中的吞吐瓶颈与消费积压问题。一、如何获取主题统计信息Pulsar 的主题统计可以通过三种方式获取命令行工具pulsar-admin、REST API 和 Java Admin API。三种方式返回的数据结构完全一致均来自 Broker 端同一条统计链路。1.1 命令行pulsar-admin统计命令统一挂在topics子命令下2.2.1 版本同时兼容旧的persistent子命令见 admin-api-persistent-topics.md$ pulsar-admin topics stats persistent://test-tenant/ns1/tp1获取主题内部统计internal stats$ pulsar-admin topics stats-internal persistent://test-tenant/ns1/tp1根据 reference-pulsar-admin.mdtopics stats命令返回主题及其连接的生产者、消费者统计且所有速率指标都在 1 分钟窗口内计算相对于最近一个完整完成的 1 分钟周期。这意味着你看到的msgRateIn/msgRateOut并非瞬时值而是上一完整分钟的平均速率解读监控曲线时应留意这一点。1.2 REST API对应 REST 端点Admin API v2普通统计GET /admin/v2/:schema/:tenant/:namespace/:topic/stats内部统计GET /admin/v2/:schema/:tenant/:namespace/:topic/internalStats1.3 Java Admin APIString topic persistent://my-tenant/my-namespace/my-topic; admin.persistentTopics().getStats(topic); // 普通统计 admin.persistentTopics().getInternalStats(topic); // 内部统计在 Broker 内部这两个端点最终都汇聚到PersistentTopicsBase的internalGetStats/internalGetInternalStats方法见 PersistentTopicsBase.java先校验全局命名空间所有权与主题操作权限TopicOperation.GET_STATS再通过topic.asyncGetStats(...)异步获取统计对象并返回。二、分区主题Partitioned Topics统计字段全解分区主题是一个逻辑主题底层由多个物理分区partition组成。调用stats命令时Broker 会聚合所有分区的统计结果。在源码层面这由 PartitionedTopicStatsImpl.java 体现它继承自TopicStatsImpl并额外持有metadata分区元数据和partitions分区名 → 各分区统计的映射。下表完整列出分区主题统计中的全部字段StatDescriptionmsgRateIn所有本地与复制replication生产者发布速率之和单位消息/秒msgThroughputIn同 msgRateIn但单位是字节/秒msgRateOut所有本地与复制消费者分发速率之和单位消息/秒msgThroughputOut同 msgRateOut但单位是字节/秒averageMsgSize最近一个统计间隔内该发布者的平均消息大小字节storageSize该主题所有 ledger 的存储大小之和字节publishers该主题上所有本地发布者的列表数量可以是 0 到数千producerId该生产者在主题上的内部标识符producerName由客户端库生成的内部生产者标识符address该生产者连接的 IP 地址与源端口connectedSince该生产者创建或最后一次重连的时间戳subscriptions该主题上所有本地订阅的列表my-subscription订阅名称由客户端定义msgBacklog该订阅的积压消息数量type订阅类型msgRateExpired该订阅因 TTL 被丢弃而非分发消息的速率consumers该订阅下已连接的消费者列表consumerName由客户端库生成的内部消费者标识符availablePermits该消费者在客户端库监听队列中还可接收的消息数。值为 0 表示客户端队列已满且未调用 receive()非 0 表示该消费者可被继续分发消息replication该主题跨机房复制的统计部分replicationBacklog出站复制积压消息数connected出站复制器outbound replicator是否已连接replicationDelayInSeconds若 connected 为 true表示最老消息在连接中等待发送的时长秒inboundConnection远端集群发布连接所连 Broker 的 IP 与端口inboundConnectedSince用于向远端集群发布消息的 TCP 连接建立时间。若没有本地发布者连接该连接会在 1 分钟后自动关闭2.1 聚合语义分区统计如何合并从源码看分区统计并不是简单“拼在一起”而是逐字段累加合并。TopicStatsImpl.add()见 TopicStatsImpl.java实现了这套合并逻辑值得注意的细节有速率与计数直接累加msgRateIn、msgThroughputIn、bytesInCounter、msgInCounter、storageSize、backlogSize等全部做加法平均消息大小按加权平均averageMsgSize使用(旧平均值 × (count-1) 新平均值) / count的递推公式随分区数量增长平滑过渡生产者按名称合并支持部分生产者partial producer场景下相同producerName的统计会合入同一个publishersMap条目并调用PublisherStatsImpl.add()累加见 PublisherStatsImpl.java避免同一客户端在多分区发布时统计条目爆炸订阅以名称对齐每个订阅名称只保留一个聚合后的SubscriptionStatsImpl复制统计按远端集群合并replicationDelayInSeconds取所有远端的最大值connected采用“与”逻辑任一断开即为断开见 ReplicatorStatsImpl.java。此外所有统计对象都提供reset()方法TopicStatsImpl.reset()、PartitionedTopicStatsImpl.reset()还会清空partitions并将metadata.partitions置 0Broker 在每次统计周期都会重置后重新采样保证返回的是当轮数据。2.2 订阅层与消费者层的关键指标解读subscriptions下的每个订阅对应 SubscriptionStatsImpl.java除了上表的msgBacklog、type、msgRateExpired之外源码中还包含更多字段可辅助排障例如unackedMessages未确认消息数、msgDelayed延迟消息数、blockedSubscriptionOnUnackedMsgs订阅是否因未确认消息达到阈值而被阻塞、activeConsumerNameExclusive/Failover 模式下的活动消费者名等。消费者的availablePermits是判断消费是否跟得上的核心指标值为 0 时说明客户端接收队列已满消费者没有主动拉取通常意味着消费端处理速度不足或未调用receive()持续非 0 且较大则表示消费侧健康、随时可接收新消息。2.3 跨机房复制的监控要点replication段来自每个复制器Replicator的实时状态。重点关注三个联动指标connected出站复制器是否存活。复制依赖内嵌的“复制订阅者”replication-subscriber把消息转发到远端集群见ReplicatorStatsImpl中msgRateOut的注释“delivered to the replication-subscriber”。replicationBacklog待复制消息数持续增长意味着本集群到远端集群的链路拥塞。replicationDelayInSeconds最老一条消息已等待复制的秒数直观反映端到端复制延迟。inboundConnectedSince若没有本地发布者向远端发布此 TCP 连接会在 1 分钟后自动关闭因此该时间戳会周期性更新属正常现象。三、主题内部统计Internal Stats / Topics字段全解stats-internal返回的是底层 BookKeeper ledger 与 cursor 级别的精细指标是深入诊断存储与消费位置的核心数据。下表为完整字段说明StatDescriptionentriesAddedCounter自该 Broker 加载该主题以来发布的消息数numberOfEntries正在跟踪的消息总数totalSize所有消息的总存储大小字节currentLedgerEntries当前正在写入的 ledger 中已写入的消息数currentLedgerSize当前正在写入的 ledger 中已写入消息的大小字节lastLedgerCreatedTimestamp最后一个 ledger 的创建时间lastLedgerCreationFailureTimestamp最后一个 ledger 创建失败的时间waitingCursorsCount有多少 cursor 已追平并等待新消息发布pendingAddEntriesCount有多少消息的异步写入请求仍在等待完成lastConfirmedEntry最后一条成功写入消息的 ledgerid:entryid。若 entryid 为 -1表示 ledger 已打开或正在打开但尚无条目写入statecursor ledger 的状态。Open 表示存在用于保存 markDeletePosition 更新的 cursor ledgerledgers该主题存放消息的所有 ledger 的有序列表cursors该主题上所有 cursor 的列表主题统计中的每个订阅都对应一个 cursormarkDeletePosition确认位置订阅者最后确认收到的消息readPosition订阅者读取消息的最新位置waitingReadOp当订阅已读到主题最新消息并等待新消息发布时为 truependingReadOps正在进行中的对 BookKeeper 的未完成读取请求计数messagesConsumedCounter自该 Broker 加载该主题以来该 cursor 确认的消息数cursorLedger用于持久化存储当前 markDeletePosition 的 ledgercursorLedgerLastEntry用于持久化存储当前 markDeletePosition 的最后一个 entryidindividuallyDeletedMessages若消息乱序确认显示 markDeletePosition 与 readPosition 之间已确认消息的范围lastLedgerSwitchTimestampcursor ledger 最近一次滚动的时间3.1 数据模型与代码出处stats-internal的返回结构对应pulsar-client-admin-api模块中的两个类ManagedLedgerInternalStats.java定义entriesAddedCounter、numberOfEntries、totalSize、currentLedgerEntries、currentLedgerSize、lastLedgerCreatedTimestamp、lastLedgerCreationFailureTimestamp、waitingCursorsCount、pendingAddEntriesCount、lastConfirmedEntry、state、ledgers、cursors等顶层字段其中LedgerInfo含ledgerId/entries/sizeCursorStats含markDeletePosition/readPosition/waitingReadOp/pendingReadOps/messagesConsumedCounter/cursorLedger/cursorLedgerLastEntry/individuallyDeletedMessages/lastLedgerSwitchTimestamp/state等字段PersistentTopicInternalStats.java在 ManagedLedger 基础上扩展schemaLedgersSchema 存储 ledger 列表与compactedLedgerTopic Compaction 产生的 ledger若存在。值得注意的是文档与示例中“state 字段”在 CursorStats 场景下指的是 cursor ledger 状态Open 表示正在保存 markDeletePosition而在主题级场景下指 ledger 写入状态LedgerOpened 表示有 ledger 处于打开状态用于保存发布的消息两种语境含义不同阅读输出时需结合上下文。3.2 一份真实的 internal stats JSON 示例以下为 admin-api-persistent-topics.md 中给出的完整返回示例可对照上表逐字段理解{ entriesAddedCounter: 20449518, numberOfEntries: 3233, totalSize: 331482, currentLedgerEntries: 3233, currentLedgerSize: 331482, lastLedgerCreatedTimestamp: 2016-06-29 03:00:23.825, lastLedgerCreationFailureTimestamp: null, waitingCursorsCount: 1, pendingAddEntriesCount: 0, lastConfirmedEntry: 324711539:3232, state: LedgerOpened, ledgers: [ { ledgerId: 324711539, entries: 0, size: 0 } ], cursors: { my-subscription: { markDeletePosition: 324711539:3133, readPosition: 324711539:3233, waitingReadOp: true, pendingReadOps: 0, messagesConsumedCounter: 20449501, cursorLedger: 324702104, cursorLedgerLastEntry: 21, individuallyDeletedMessages: [(324711539:3134‥324711539:3136], (324711539:3137‥324711539:3140], ], lastLedgerSwitchTimestamp: 2016-06-29 01:30:19.313, state: Open } ] }3.3 如何读懂这份输出markDeletePosition与readPosition324711539:3133vs324711539:3233两者差距即该订阅尚未确认消费的消息区间。若差值长期扩大而messagesConsumedCounter增速放缓说明消费端处理能力不足若waitingReadOp为 true说明该订阅已经追平到最新消息正在等待新消息属于健康状态。pendingReadOps/pendingAddEntriesCount这两个值应长期接近 0。若持续偏高意味着对 BookKeeper 的异步读写请求大量堆积通常是存储集群响应变慢或网络抖动的前兆。individuallyDeletedMessages乱序确认产生的空洞区间。(324711539:3134‥324711539:3136]表示已被确认但位于 markDeletePosition 之后的若干条消息。区间过多会增加消费时的跳跃成本也侧面反映业务侧乱序 ack 的频率。cursorLedger与lastLedgerSwitchTimestampcursor 位置持久化在专用 ledger 中周期滚动rollover频繁滚动通常对应高频 ack属预期行为。lastLedgerCreationFailureTimestamp若此字段非 null说明曾有 ledger 创建失败需要结合 Broker 日志与 BookKeeper 集群状态排查存储异常。四、从源码理解统计的生命周期总结起来Pulsar 统计在源码中遵循清晰的“三件套”模式以TopicStatsImpl为代表见 TopicStatsImpl.java采样Broker 的持久化主题PersistentTopic在消息写入、分发、确认、复制等关键路径上实时更新各计数字段聚合分区主题通过add()将各分区统计合并为一份逻辑统计内部count字段记录参与聚合的分区数用于加权平均等计算且被JsonIgnore标记不对外输出输出通过 Admin API /pulsar-admin/ Java Admin API 以 JSON 形式暴露速率类字段按 1 分钟窗口计算。理解这三层之后你在pulsar-admin topics stats与stats-internal之间切换时就能做到心中有数前者回答“业务侧快不快”速率、积压、消费者状态后者回答“存储侧稳不稳”ledger、cursor、读写堆积两者配合即可对单个主题的完整链路做一次系统体检。本文基于仓库 version-2.2.1 版 administration-stats 文档 撰写并结合 admin-api-persistent-topics.md 与 reference-pulsar-admin.md 中的命令示例以及pulsar-common/pulsar-client-admin-api/pulsar-broker模块源码进行佐证。如需查看当前版本各字段的完整 JSON 输出可直接运行pulsar-admin topics stats与pulsar-admin topics stats-internal命令。赞分享消息队列后端流处理【免费下载链接】pulsarApache Pulsar - distributed pub-sub messaging system项目地址https://gitcode.com/gh_mirrors/pulsar28/pulsar点击查看免费下载相关推荐Apache Pulsar 主题统计Topic Stats字段详解从 Admin API 到内部存储的完整解读Apache Pulsar 主题统计Topic Stats字段详解从 Admin API 到内部存储的完整解读 导读 Pulsar 的统计stats体消息队列后端流处理Apache Pulsar 主题统计指标全解析分区主题 stats 与内部 stats 字段详解Apache Pulsar 主题统计指标全解析分区主题 stats 与内部 stats 字段详解 本指南以 Apache Pulsar 的 administr消息队列后端流处理Apache Pulsar 统计指标详解Partitioned Topics 与 Topic 内部统计字段全解读Apache Pulsar 统计指标详解Partitioned Topics 与 Topic 内部统计字段全解读 Apache Pulsar 内置了完整的运行消息队列后端流处理创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考
返回列表