- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
本文基于 Apache Pulsar 仓库中 deploy-monitoring.md(version-2.2.0 文档)整理,并结合当前仓库的源码、配置文件与 Grafana 仪表盘模板加以印证。文章面向运维与开发人员,系统讲解 Pulsar 集群中 Broker、ZooKeeper、BookKeeper、Functions/Connectors 四大类组件的指标采集方式,以及如何接入 Prometheus 抓取、使用 Grafana 展示,并配置告警规则,帮助读者搭建一套可观测、可告警的生产级监控体系。
监控体系总览:能监控什么、用什么方式监控
Pulsar 集群由多个组件构成,因此监控也需要分层进行。官方文档明确指出,一个 Pulsar 集群的监控既要覆盖主题(topic)维度的使用量指标,也要覆盖集群各组件的整体健康状态。
从指标来源上看,监控数据分布在四类组件中:
| 组件 | 指标内容 | 采集方式 |
|---|---|---|
| Broker | 主题级统计(destination dumps)、命名空间级聚合指标 | pulsar-admin broker-stats命令 + Prometheus HTTP 端点 |
| ZooKeeper | 本地 ZooKeeper、配置存储(configuration store)服务器与客户端的详细统计 | Prometheus HTTP 端点(/metrics) |
| BookKeeper | bookie 的统计指标 | conf/bookkeeper.conf中配置的 stats 框架(默认 Prometheus 导出器) |
| Functions/Connectors | functions-worker 的 JVM 指标、函数与连接器指标 | pulsar-admin functions-worker命令 + Prometheus HTTP 端点 |
下面按组件逐一说明指标采集的具体方式与底层实现。
采集 Broker 指标
两类 Broker 统计:destination dumps 与 monitoring metrics
Pulsar Broker 的指标以 JSON 格式对外暴露,主要分为两类:
第一类:Destination dumps(主题级统计)。包含每个独立主题(topic)的统计信息,使用如下命令获取:
bin/pulsar-admin broker-stats destinations第二类:Broker metrics(命名空间级聚合指标)。包含 Broker 自身信息以及按命名空间(namespace)聚合的主题统计,使用如下命令获取:
bin/pulsar-admin broker-stats monitoring-metrics注意:所有消息速率(message rates)指标每1 分钟更新一次。
从源码可以看到这两类指标的实现细节。在 CmdBrokerStats.java 中,broker-stats命令注册了monitoring-metrics、mbeans、topics(别名destinations)、allocator-stats、load-report等子命令:
public CmdBrokerStats(Supplier<PulsarAdmin> admin) { super("broker-stats", admin); jcommander.addCommand("monitoring-metrics", new CmdMonitoringMetrics()); jcommander.addCommand("mbeans", new CmdDumpMBeans()); jcommander.addCommand("topics", new CmdTopics(), "destinations"); jcommander.addCommand("allocator-stats", new CmdAllocatorStats()); jcommander.addCommand("load-report", new CmdLoadReport()); }也就是说,destinations实际上是topics子命令的别名,其实现调用getAdmin().brokerStats().getTopics(),将每个主题的统计以 JSON 输出;monitoring-metrics则调用getAdmin().brokerStats().getMetrics(),支持-i/--indent参数对 JSON 输出进行缩进美化。
在 REST 服务端,BrokerStatsBase.java 将/metrics路径映射为监控采集端点,其 API 注解明确说明"该请求应由监控代理在每个 broker 上执行以抓取指标"("Requested should be executed by Monitoring agent on each broker to fetch the metrics"),并且仅允许超级用户(super user)访问:
@GET @Path("/metrics") public Collection<Metrics> getMetrics() throws Exception { // Ensure super user access only validateSuperUserAccess(); Collection<Metrics> metrics = pulsar().getMetricsGenerator().generate(); return metrics; }Prometheus 格式的聚合指标
除了 JSON 格式之外,聚合后的 Broker 指标还会以 Prometheus 格式在如下地址暴露:
http://$BROKER_ADDRESS:8080/metrics/其中8080是 Broker 的默认 Web 服务端口,$BROKER_ADDRESS替换为具体 broker 的主机地址。该端点供 Prometheus 直接抓取(scrape),无需额外开发导出器。
采集 ZooKeeper 指标
Pulsar 自带的本地 ZooKeeper、配置存储(configuration store)服务器及其客户端,都可以通过 Prometheus 暴露详细的统计指标:
http://$LOCAL_ZK_SERVER:8000/metrics http://$GLOBAL_ZK_SERVER:8001/metrics其中:
- 本地 ZooKeeper 的默认端口为8000;
- 配置存储(全局 ZooKeeper)的默认端口为8001。
如果要修改这两个默认端口,可以通过指定系统属性(system property)stats_server_port来完成。从架构上看,Pulsar 的元数据服务(pulsar-metadata模块)负责 ZooKeeper 统计端口的启动与监听,修改该属性后重启相关服务即可生效。
采集 BookKeeper 指标
BookKeeper 的统计框架通过修改 conf/bookkeeper.conf 中的statsProviderClass来配置。当前仓库默认配置如下(conf/bookkeeper.conf):
statsProviderClass=org.apache.bookkeeper.stats.prometheus.PrometheusMetricsProvider prometheusStatsHttpPort=8000也就是说:
- 默认的 BookKeeper 配置已经启用 Prometheus 导出器(
PrometheusMetricsProvider),该配置随 Pulsar 发行包一起分发; - bookie 的 Prometheus 指标地址为:
http://$BOOKIE_ADDRESS:8000/metrics- bookie 指标端口的默认值为8000,可通过
conf/bookkeeper.conf中的prometheusStatsHttpPort修改。
跟踪 Managed Cursor 确认状态的指标
在 Pulsar 中,消费者确认(acknowledgment)状态会优先持久化到 ledger;当写入 ledger 失败时,才会回退持久化到 ZooKeeper。为了跟踪确认过程的统计信息,可以为 Managed Cursor 配置如下 Prometheus 指标:
brk_ml_cursor_persistLedgerSucceed(namespace="", ledger_name="", cursor_name:"") brk_ml_cursor_persistLedgerErrors(namespace="", ledger_name="", cursor_name:"") brk_ml_cursor_persistZookeeperSucceed(namespace="", ledger_name="", cursor_name:"") brk_ml_cursor_persistZookeeperErrors(namespace="", ledger_name="", cursor_name:"") brk_ml_cursor_nonContiguousDeletedMessagesRange(namespace="", ledger_name="", cursor_name:"")这些指标以namespace、ledger_name、cursor_name为标签维度,含义如下:
| 指标 | 含义 |
|---|---|
brk_ml_cursor_persistLedgerSucceed | 确认状态成功持久化到 ledger 的次数 |
brk_ml_cursor_persistLedgerErrors | 确认状态持久化到 ledger 时失败的次数 |
brk_ml_cursor_persistZookeeperSucceed | 确认状态成功持久化到 ZooKeeper 的次数 |
brk_ml_cursor_persistZookeeperErrors | 确认状态持久化到 ZooKeeper 时失败的次数 |
brk_ml_cursor_nonContiguousDeletedMessagesRange | 非连续删除消息区间(与消息删除/游标推进相关) |
这些指标会添加到 Prometheus 接口中,可在 Grafana 中监控和查看。源码层面,ManagedCursorMXBeanImpl.java 维护了persistZookeeperSucceed等 LongAdder 计数器,并通过 MBean 方式暴露给统计系统,正是上述指标的底层数据来源。
采集 Functions 与 Connectors 指标
Pulsar Functions 与 Connectors 运行于 functions-worker 进程中,其指标同样支持 JSON 与 Prometheus 两种格式。
JSON 格式
functions-worker JVM 指标(包含 functions worker 的 JVM 指标),使用命令:
pulsar-admin functions-worker monitoring-metrics函数与连接器指标,使用命令:
pulsar-admin functions-worker function-stats从 CmdFunctionWorker.java 可以看到,functions-worker命令下注册了function-stats与monitoring-metrics两个子命令,与官方文档一一对应。底层实现中,MetricsGenerator.java 将 JVM 指标等聚合为Metrics列表返回;而 FunctionsStatsGenerator.java 负责从各个函数运行时实例收集 Prometheus 格式的指标文本。
Prometheus 格式
聚合后的函数与连接器指标以 Prometheus 格式暴露在如下地址:
http://$FUNCTIONS_WORKER_ADDRESS:$WORKER_PORT/metrics其中$FUNCTIONS_WORKER_ADDRESS(functions worker 地址)和$WORKER_PORT(worker 端口)都来自 conf/functions_worker.yml 配置文件。配置文件中对应字段为workerHostname(或workerId中解析出的地址)与workerPort,部署时按实际环境替换即可。
配置 Prometheus 抓取指标
收集到各组件暴露的指标后,下一步就是让 Prometheus 定期抓取。官方建议的工作方式是:用 Prometheus 采集 Pulsar 各组件暴露的全部指标,再用 Grafana 仪表盘展示并监控集群状态。
- 裸机(bare metal)部署:需要手动在 Prometheus 配置中提供要探测的节点列表,即把上述各组件的指标端点加入
scrape_configs,例如:
scrape_configs: - job_name: 'pulsar-broker' static_configs: - targets: ['$BROKER_ADDRESS:8080'] metrics_path: /metrics - job_name: 'pulsar-bookie' static_configs: - targets: ['$BOOKIE_ADDRESS:8000'] metrics_path: /metrics - job_name: 'pulsar-zookeeper' static_configs: - targets: - '$LOCAL_ZK_SERVER:8000' - '$GLOBAL_ZK_SERVER:8001' metrics_path: /metrics- Kubernetes 部署:监控会自动配置,无需手工维护抓取列表(详见 deploy-kubernetes.md 中的 Kubernetes 部署说明)。
使用 Grafana 仪表盘
维度控制的设计原则
当开始采集时间序列数据时,最大的挑战是防止数据附加的维度(dimension)数量爆炸。因此官方文档强调:在时间序列采集层面,只需要采集按命名空间聚合的指标即可满足大多数监控需求,避免为每个主题生成过多的高基数时间序列导致存储与查询压力过大。
Grafana 仪表盘部署方式
Grafana 可以直接使用 Prometheus 中存储的数据创建仪表盘。在 Kubernetes 上部署 Pulsar 时,默认会启用pulsar-grafanaDocker 镜像,该镜像自带主要的仪表盘。如需手动启动该镜像,使用如下命令:
docker run -p3000:3000 \ -e PROMETHEUS_URL=http://$PROMETHEUS_HOST:9090/ \ apachepulsar/pulsar-grafana:latest参数说明:
-p3000:3000:将容器内 Grafana 的 3000 端口映射到宿主机;-e PROMETHEUS_URL=...:通过环境变量指定 Prometheus 服务地址,$PROMETHEUS_HOST替换为实际 Prometheus 主机,9090是 Prometheus 默认端口。
仓库内置的 Grafana 仪表盘模板
当前仓库的 grafana/dashboards 目录内置了多份可直接导入 Grafana 的仪表盘 JSON 模板,覆盖集群各核心组件:
| 文件 | 覆盖组件/视角 |
|---|---|
| broker.json | Broker 运行指标 |
| bookkeeper.json | BookKeeper 指标 |
| zookeeper.json | ZooKeeper 指标 |
| namespace.json | 命名空间级聚合指标 |
| topic.json | 主题级指标 |
| jvm.json | 各组件 JVM 指标 |
| prometheus.json | Prometheus 自身指标 |
这些模板与pulsar-grafana镜像配合使用,可以直接在 Grafana 中 Import 加载。此外,Pulsar Manager 也提供了逐主题(per-topic)仪表盘的说明,可参考 administration-dashboard.md。
配置告警规则
监控的最终目的是及时发现问题。官方文档建议根据自身的 Pulsar 环境设置告警规则,典型思路包括:
- 对
brk_ml_cursor_persistLedgerErrors等错误类计数器设置阈值告警,当持续增长时触发; - 对 broker/bookie 端口(如
8080、8000)的可达性设置探活告警; - 对命名空间级消息速率、积压(backlog)等指标设置容量告警。
Prometheus 的告警规则通过rules文件配置,例如:
groups: - name: pulsar-alerts rules: - alert: CursorPersistLedgerErrors expr: increase(brk_ml_cursor_persistLedgerErrors[5m]) > 0 for: 5m labels: severity: warning annotations: summary: "Managed cursor ledger persist errors increasing"具体规则语法以 Prometheus 告警规则文档为准(Prometheus Alerting rules),配置完成后由 Alertmanager 负责发送通知。
小结
一个完整的 Pulsar 监控体系可以归纳为三条链路:
- 指标暴露:Broker(JSON +
:8080/metrics)、ZooKeeper(:8000/:8001)、BookKeeper(默认:8000)、Functions/Connectors(functions_worker.yml指定的:WORKER_PORT)各自暴露 Prometheus 格式指标; - 指标采集与展示:Prometheus 定期抓取各端点,Grafana 通过内置模板(见 grafana/dashboards)或
pulsar-grafana镜像可视化展示; - 告警闭环:依据环境特征配置 Prometheus 告警规则,并结合 managed cursor 确认状态、命名空间级聚合指标等关键信号及时发现问题。
部署方式上,裸机环境需手工维护抓取节点列表,Kubernetes 环境则由部署组件自动配置监控。读者可结合本文列出的源码路径(如 CmdBrokerStats.java、BrokerStatsBase.java、conf/bookkeeper.conf)进一步深入理解指标的产生与暴露机制,从而按需定制自己的监控方案。
- 消息队列
- 后端
- 流处理
【免费下载链接】pulsar
Apache Pulsar - distributed pub-sub messaging system
相关推荐
Apache Pulsar 集群监控部署指南:指标采集、Prometheus 配置与 Grafana 可视化
Apache Pulsar 集群监控部署指南:指标采集、Prometheus 配置与 Grafana 可视化 Apache Pulsar 是一个分布式 pub
消息队列后端流处理Apache Pulsar 集群监控部署指南:指标采集、Prometheus 配置与 Grafana 面板实践
Apache Pulsar 集群监控部署指南:指标采集、Prometheus 配置与 Grafana 面板实践 本篇技术指南围绕 Apache Pulsar 集
消息队列后端流处理Apache Pulsar 集群监控实践:Metrics 采集、Prometheus 对接与 Grafana 仪表盘
Apache Pulsar 集群监控实践:Metrics 采集、Prometheus 对接与 Grafana 仪表盘 Apache Pulsar 提供了多层次的
消息队列后端流处理
创作声明:本文部分内容由AI辅助生成(AIGC),仅供参考