第1章:Pulsar 概述与核心概念
1.1 什么是消息系统与流处理平台
| 概念名称 | 说明 | 注意事项 |
|---|
| 消息系统(Message System) | 用于在分布式系统组件之间传递数据的中间件,实现解耦、异步通信和流量削峰。典型模式包括点对点(Queue)和发布-订阅(Pub-Sub)。 | 需关注消息的可靠性(是否丢失)、顺序性、吞吐量和延迟。 |
| 流处理平台(Streaming Platform) | 支持实时数据流的采集、存储、处理与分发的系统,不仅传输消息,还支持对数据流进行计算(如聚合、过滤)。 | 区别于批处理,强调低延迟和持续处理能力。 |
| 消息(Message) | 数据传输的基本单位,通常包含 payload(数据体)、key(路由键)、properties(元数据)、timestamp 等字段。 | 应控制单条消息大小(通常建议 ≤ 5MB),避免网络阻塞。 |
| 生产者(Producer) | 向消息系统发送消息的客户端程序。 | 需处理发送失败、超时、重试等异常情况。 |
| 消费者(Consumer) | 从消息系统接收并处理消息的客户端程序。 | 需正确管理消费位点(cursor),避免重复或丢失消息。 |
| 主题(Topic) | 消息的逻辑分类通道,生产者发送消息到 Topic,消费者从 Topic 订阅消息。 | 一个 Topic 可被多个消费者组消费,实现广播或负载均衡。 |
1.2 Pulsar 的诞生背景与架构演进
| 概念名称 | 说明 | 注意事项 |
|---|
| 传统消息系统瓶颈 | Kafka 等系统将存储与计算耦合,扩展性受限,运维复杂,难以支持多租户和云原生环境。 | 耦合架构导致扩容需整体迁移,影响可用性。 |
| Pulsar 起源 | 由 Yahoo 开发,2016 年开源,2018 年成为 Apache 顶级项目,旨在解决大规模、多租户、高可用的消息场景。 | 初始设计目标是支持 Yahoo 内部数十万个 Topic 和百万级 QPS。 |
| 分层架构思想 | 将消息服务层(无状态 Broker)与存储层(BookKeeper)分离,实现计算与存储独立扩展。 | 存储层专注于数据持久化,服务层专注于请求路由与协议处理。 |
| 云原生支持 | 天然支持 Kubernetes 部署,组件可独立扩缩容,适合微服务架构。 | 通过 Helm Chart 可快速部署 Pulsar 集群。 |
| 架构演进关键点 | 早期:单体架构 → 中期:引入 BookKeeper 实现分层 → 现代:支持 Functions、IO、分层存储、Geo-Replication | 演进方向是功能丰富化与云原生深度集成。 |
1.3 Pulsar 与 Kafka 的核心对比
| 对比维度 | Apache Pulsar | Apache Kafka | 注意事项 |
|---|
| 架构模式 | 分层架构(Broker + BookKeeper) | 存算耦合(Broker 自带存储) | Pulsar 更易独立扩展存储或计算资源。 |
| 存储引擎 | Apache BookKeeper(分布式日志) | 自研基于文件系统的日志存储 | BookKeeper 提供更低的尾延迟和更高的写入可用性。 |
| 多租户支持 | 原生支持,通过命名空间(Namespace)隔离资源 | 需额外工具或配置实现,非原生 | Pulsar 更适合 SaaS、企业级平台。 |
| 订阅模式 | 支持 Exclusive、Shared、Failover、Key_Shared | 仅支持 Consumer Group(类似 Shared) | Pulsar 的 Key_Shared 模式支持按 Key 路由,避免重复消费。 |
| 消息确认机制 | 支持 Individual Ack 和 Cumulative Ack | 仅支持 Offset 提交(Cumulative) | Pulsar 可精确确认单条消息,避免重处理。 |
| 延迟消息 | 原生支持延迟投递 | 需依赖外部系统(如时间轮)或客户端实现 | Pulsar 使用分层存储或定时器实现延迟消息。 |
| 分层存储 | 支持将冷数据卸载到 S3、GCS 等对象存储 | 社区版不支持,商业版支持 | 降低长期存储成本,适合日志归档场景。 |
| 运维复杂度 | 组件多(ZooKeeper、BookKeeper、Broker),初始部署较复杂 | 相对简单,但扩容需数据迁移 | Pulsar 运维学习曲线较高,但长期扩展性更好。 |
1.4 Pulsar 的四大核心优势
| 优势 | 说明 | 注意事项 |
|---|
| 分层架构 | 计算(Broker)与存储(BookKeeper)分离,Broker 无状态,可快速扩缩容;BookKeeper 专注高可用存储。 | 需维护多个组件,但故障隔离性更好。 |
| 多租户支持 | 通过命名空间(Namespace)实现租户隔离,支持配额、认证、ACL、复制策略等。 | 适合企业级平台或云服务提供商。 |
| 持久化与高可用 | 所有消息默认持久化到 BookKeeper,支持多副本(ensemble)和自动故障恢复。 | 可配置 ackQuorum 和 writeQuorum 保证写入一致性。 |
| 低延迟与高吞吐 | 即使在高 backlog 场景下,仍能保持低尾延迟(Tail Latency),适合实时场景。 | 得益于 BookKeeper 的分片写入和读写分离机制。 |
1.5 Pulsar 核心组件概览
| 组件 | 说明 | 注意事项 |
|---|
| Broker | 无状态服务节点,负责接收生产者消息、推送消息给消费者、管理 Topic 路由。 | 不存储消息,宕机后可由其他 Broker 接管,不影响数据。 |
| BookKeeper | 分布式持久化日志存储系统,由多个 Bookie 组成,存储消息数据。 | 需部署奇数个 Bookie(如 3、5)以实现多数派写入。 |
| ZooKeeper | 存储集群元数据(如 Broker 注册、Topic 配置、租户策略),不参与消息传输。 | 高可用部署(至少 3 节点),是集群”大脑”。 |
| Proxy | 可选组件,用于统一接入层,支持 TLS 终止、身份验证、负载均衡。 | 在 Kubernetes 或公有云中常用于暴露服务入口。 |
| Configuration Store | 通常为独立的 ZooKeeper 集群,存储租户、命名空间等配置信息。 | 与本地 ZooKeeper 分离,避免配置与运行时数据竞争。 |
第2章:Pulsar 架构深度解析
2.1 整体架构图解与数据流路径
| 组件/路径 | 说明 | 注意事项 |
|---|
| 生产者 → Broker | 生产者连接 Broker,发送消息;Broker 将消息转发给 BookKeeper。 | Broker 不缓存消息,直接写入 BookKeeper。 |
| Broker → BookKeeper | Broker 将消息封装为 Entry,写入 BookKeeper 的 Ledger。 | 写入是异步的,支持批处理提升吞吐。 |
| BookKeeper 存储 | 多个 Bookie 组成 Ensemble,消息被分片并多副本存储。 | 默认 ackQuorum=2,writeQuorum=2,ensemble=3。 |
| 消费者 ← Broker | 消费者从 Broker 拉取消息,Broker 从 BookKeeper 读取数据并推送。 | 读取路径:Broker → Bookie → Consumer。 |
| 元数据 ← ZooKeeper | Broker 启动时从 ZooKeeper 获取 Topic 路由、租户配置等信息。 | ZooKeeper 不参与消息读写,仅管理元数据。 |
| 数据流总结 | Producer → Broker → BookKeeper(持久化),Consumer ← Broker ← BookKeeper(读取) | 数据流清晰分离,写入与读取路径解耦。 |
2.2 Broker 的职责与负载均衡机制
| 功能 | 说明 | 注意事项 |
|---|
| 客户端连接管理 | 接收 Producer 和 Consumer 的连接请求,维护会话状态。 | 支持多种协议(Pulsar、Kafka via Kafka-on-Pulsar)。 |
| Topic 路由 | 根据 Topic 名称查找归属的 Bundle,确定由哪个 Broker 服务。 | 使用一致性哈希(load balance bundle)分配 Topic。 |
| 消息转发 | 将生产者消息转发给 BookKeeper,将 BookKeeper 数据推送给消费者。 | 不存储消息,是”中转站”。 |
| 负载均衡策略 | 支持动态和静态两种模式:Dynamic:基于 CPU、内存、带宽等指标自动迁移 Topic;Static:固定分配。 | 建议生产环境使用 Dynamic,避免热点。 |
| 故障转移 | Broker 宕机后,ZooKeeper 检测到失联,其他 Broker 接管其 Topic。 | Topic 所有者变更不影响数据,因存储在 BookKeeper。 |
| 协议支持 | 原生 Pulsar 协议,可通过 Kafka-on-Pulsar 支持 Kafka 客户端。 | Kafka 兼容层性能略低于原生协议。 |
2.3 BookKeeper:分布式日志存储原理
| 概念 | 说明 | 注意事项 |
|---|
| Bookie | BookKeeper 的存储节点,负责存储 Entry(消息条目)。 | 每个 Bookie 是独立的存储单元,需 SSD 提升性能。 |
| Ledger | 逻辑日志单元,Pulsar 的每个 Topic 分区对应一个或多个 Ledger。 | Ledger 不可变,写满后自动滚动创建新 Ledger。 |
| Entry | 消息在 BookKeeper 中的存储单位,包含消息数据和元数据。 | Entry ID 递增,保证顺序写入。 |
| Ensemble | 参与写入的 Bookie 集合,如 ensemble=3 表示 3 个 Bookie。 | 写入时消息被复制到 ensemble 中的多个 Bookie。 |
| Write Quorum | 必须成功写入的 Bookie 数量(如 2/3)。 | writeQuorum ≤ ensembleSize。 |
| Ack Quorum | 返回成功前必须确认写入的 Bookie 数量(如 2/3)。 | ackQuorum ≤ writeQuorum,保证强一致性。 |
| 分片写入 | 消息被分片并行写入多个 Bookie,提升吞吐和降低延迟。 | 即使部分 Bookie 慢,也不阻塞整体写入。 |
2.4 ZooKeeper 的元数据管理作用
| 元数据类型 | 说明 | 注意事项 |
|---|
| Broker 注册 | Broker 启动时在 ZooKeeper 注册,形成集群视图。 | ZooKeeper 监控 Broker 存活性(心跳机制)。 |
| Topic 所有者信息 | 记录每个 Topic 当前由哪个 Broker 服务。 | Broker 宕机后,其他 Broker 可通过 ZooKeeper 重新分配。 |
| 租户与命名空间配置 | 存储多租户的配额、策略、ACL 等信息。 | 配置变更通过 ZooKeeper 通知所有 Broker。 |
| Schema 信息 | 存储消息 Schema 定义与版本。 | 支持 Schema 演化与兼容性检查。 |
| 订阅状态(Cursor) | 存储消费者组的消费位点(位置),实现持久化跟踪。 | 即使消费者重启,也能从上次位置继续消费。 |
| 配置中心 | 作为全局配置存储,支持动态更新(如 TTL、Retention)。 | 避免重启 Broker 生效配置。 |
2.5 Topic、Partition 与 Subscription 模型
| 概念 | 说明 | 注意事项 |
|---|
| Topic | 消息的逻辑通道,格式为 persistent://租户/命名空间/主题名,例如:persistent://public/default/my-topic | 必须属于某个命名空间。 |
| Persistent Topic | 消息持久化到 BookKeeper,重启不丢失。 | 默认类型,适用于大多数场景。 |
| Non-Persistent Topic | 消息不持久化,仅在内存中传输,Broker 重启后丢失。 | 适用于低延迟、可丢失场景(如监控数据)。 |
| Partitioned Topic | 一个 Topic 拆分为多个分区(Partition),每个分区是一个独立的 Topic。 | 提升并行度和吞吐量,最大分区数需预先设置。 |
| Subscription | 消费者组的抽象,表示一组消费者共同消费一个 Topic。 | 一个 Topic 可有多个 Subscription。 |
| Subscription 类型 | Exclusive:单个消费者;Shared:多个消费者共享,无序;Failover:多个消费者,主备模式;Key_Shared:按 Key 路由,保证同 Key 顺序 | Key_Shared 需指定分发 Key。 |
2.6 消息分发模式(Publish-Subscribe、Queue、Key-Shared)
| 分发模式 | 说明 | 代码示例(Java) | 注意事项 |
|---|
| Publish-Subscribe (Pub-Sub) | 一个消息被多个订阅者接收,每个 Subscription 独立消费全量消息。 | 见下方代码块 | 适合广播场景,如通知、日志分发。 |
| Queue 模式(Shared Subscription) | 多个消费者共享一个 Subscription,消息被负载均衡到任一消费者(无序)。 | .subscriptionType(SubscriptionType.Shared) | 消费者数量可动态增减,适合任务队列。 |
| Key-Shared 模式 | 消息按 Key 哈希,相同 Key 的消息由同一消费者处理,保证顺序。 | .subscriptionType(SubscriptionType.Key_Shared).keySharedPolicy(KeySharedPolicy.stickyHash()) | 适用于需要按 Key 顺序处理的场景(如订单流)。 |
| Failover 模式 | 多个消费者注册,按优先级主备切换,主消费者失效后由备消费者接管。 | .subscriptionType(SubscriptionType.Failover) | 保证顺序性和高可用,但并发度低。 |
Pub-Sub 代码示例:
PulsarClient client = PulsarClient.builder()
.serviceUrl("pulsar://localhost:6650")
.build();
Consumer<byte[]> consumer = client.newConsumer()
.topic("my-topic")
.subscriptionName("sub1")
.subscriptionType(SubscriptionType.Shared)
.subscribe();
第3章:Pulsar 部署与运维
3.1 单机模式部署(Standalone)
| 操作步骤 | 说明 | 注意事项 |
|---|
| 下载 Pulsar | 从 Apache 官网下载最新稳定版本(如 3.3.0) | 建议选择 -bin 包,不含源码 |
| 解压安装包 | tar -xzf apache-pulsar-3.3.0-bin.tar.gz | 可移动到 /opt/pulsar 等标准路径 |
| 启动 Standalone | bin/pulsar standalone | 默认启用 Broker、BookKeeper、ZooKeeper、Proxy |
| 访问管理界面 | 浏览器打开 http://localhost:8080 | Web UI 端口为 8080,Broker 服务端口 6650 |
| 停止服务 | Ctrl+C 或 bin/pulsar-daemon stop standalone | 推荐使用 daemon 模式后台运行 |
- ✅ 适用场景:本地开发、测试、学习
- ❌ 不适用:生产环境(无高可用、性能受限)
| 部署方式 | 步骤概要 | 注意事项 |
|---|
| Docker Compose | 使用 docker-compose.yml 启动 ZooKeeper、BookKeeper、Broker、Proxy | 适合多节点测试,需配置网络和卷映射 |
| Kubernetes | 使用 Helm Chart 部署(helm install pulsar pulsar/pulsar) | 支持自动扩缩容、滚动更新,生产推荐 |
| Bare Metal(物理机) | 手动部署各组件,配置 conf/ 下的配置文件 | 控制最精细,但运维复杂度高 |
| 组件分布建议 | ZooKeeper:3/5 节点;BookKeeper:3+ 节点;Broker:2+ 节点;Proxy:边缘节点部署 | BookKeeper 建议使用 SSD,Broker 可用普通磁盘 |
- ✅ Kubernetes 部署优势:自动故障恢复、配置管理、服务发现
- 🔧 关键命令:
helm repo add pulsar https://pulsar.apache.org/charts
3.3 配置文件详解(broker.conf、bookkeeper.conf、zk.conf)
broker.conf 核心配置
| 配置项 | 语法 | 用途 | 示例值 | 注意事项 |
|---|
| brokerServiceUrl | brokerServiceUrl=pulsar://localhost:6650 | Broker 服务地址(客户端连接) | pulsar://broker1:6650 | 多节点需配置为域名或 VIP |
| webServiceUrl | webServiceUrl=http://localhost:8080 | HTTP 接口地址(Admin API) | http://broker1:8080 | 用于 REST 管理接口 |
| zookeeperServers | zookeeperServers=localhost:2181 | 连接本地 ZooKeeper 地址 | zk1:2181,zk2:2181 | 必须可达 |
| configurationStoreServers | configurationStoreServers=localhost:2181 | 配置存储 ZooKeeper 地址 | 同上 | 多集群时指向全局 ZooKeeper |
| managedLedgerDefaultEnsembleSize | managedLedgerDefaultEnsembleSize=3 | 默认写入的 Bookie 数量 | 3 | 必须 ≤ Bookie 总数 |
| managedLedgerDefaultWriteQuorum | managedLedgerDefaultWriteQuorum=2 | 写入成功需确认的副本数 | 2 | ≤ ensembleSize |
| managedLedgerDefaultAckQuorum | managedLedgerDefaultAckQuorum=2 | 返回成功前需确认的副本数 | 2 | ≤ writeQuorum |
bookkeeper.conf 核心配置
| 配置项 | 语法 | 用途 | 示例值 | 注意事项 |
|---|
| bookiePort | bookiePort=3181 | Bookie 服务端口 | 3181 | 需防火墙开放 |
| zkServers | zkServers=localhost:2181 | 连接 ZooKeeper | 同上 | 必须与 Pulsar 共享 |
| journalDirectory | journalDirectory=/pulsar/journal | 事务日志存储路径 | 自定义路径 | 建议使用独立高速磁盘 |
| ledgerDirectories | ledgerDirectories=/pulsar/ledgers | 数据文件存储路径 | 多路径可用逗号分隔 | 提升 IO 并行度 |
| useV2WireProtocol | useV2WireProtocol=true | 启用 V2 协议提升性能 | true | 建议开启 |
zk.conf 配置要点
| 配置项 | 说明 | 注意事项 |
|---|
| dataDir | ZooKeeper 数据存储目录 | 需定期清理 snapshot |
| clientPort | 客户端连接端口(默认 2181) | 多实例部署时需区分 |
| tickTime | 心跳间隔(毫秒) | 默认 2000,不建议修改 |
| initLimit | Follower 初始化同步时限 | 通常 10 |
| syncLimit | Follower 与 Leader 同步时限 | 通常 5 |
3.4 监控与指标收集(Prometheus + Grafana)
| 组件 | 指标类型 | 用途 | 配置方式 | 注意事项 |
|---|
| Pulsar Broker | JVM、Topic 数、生产/消费速率、延迟 | 监控服务健康与负载 | 在 broker.conf 中启用:prometheusStatsHttpPort=8080、exposeTopicLevelMetricsInPrometheus=true | 需配置 Prometheus scrape_configs |
| BookKeeper | Bookie IO、Ledger 写入延迟、Entry 处理速率 | 存储层性能监控 | bookkeeper.conf 中设置:statsProviderClass=org.apache.bookkeeper.stats.prometheus.PrometheusMetricsProvider | 需引入 Prometheus 依赖 |
| ZooKeeper | 请求延迟、ZNode 数、Watcher 数 | 元数据层稳定性 | 使用 zk-exporter 暴露指标 | 独立部署 exporter |
| Grafana 面板 | 可视化展示 | 快速定位瓶颈 | 导入官方 Dashboard(ID: 10000+) | 支持多集群视图 |
Prometheus 配置片段示例:
scrape_configs:
- job_name: 'pulsar'
static_configs:
- targets: ['broker1:8080', 'broker2:8080']
- job_name: 'bookkeeper'
static_configs:
- targets: ['bookie1:8080']
3.5 日志管理与故障排查
| 日志文件 | 路径 | 用途 | 常见问题 | 排查建议 |
|---|
| broker.log | logs/pulsar-broker-*.log | Broker 运行日志 | 启动失败、Topic 创建异常 | 检查 ZooKeeper 连接、端口占用 |
| bookkeeper.log | logs/pulsar-bookkeeper-*.log | Bookie 写入日志 | 写入超时、Ledger 错误 | 检查磁盘空间、journal 同步性能 |
| zookeeper.log | logs/pulsar-zookeeper-*.log | ZK 选举与事务日志 | Leader 选举失败、Session 超时 | 检查网络延迟、tickTime 配置 |
| GC 日志 | logs/gc.log | JVM 垃圾回收情况 | 频繁 Full GC、STW 过长 | 调整堆大小、使用 G1GC |
常见命令:
bin/pulsar-admin brokers list — 查看集群状态
bin/pulsar-admin topics list public/default — 列出 Topic
bin/pulsar-admin bookies list-bookies — 列出 Bookie
故障排查流程:
- 检查日志错误关键词(ERROR、Exception)
- 验证组件间网络连通性(telnet 端口)
- 使用
pulsar-admin 检查资源状态
- 查看 Prometheus 指标趋势(如延迟突增)
3.6 多租户与命名空间管理
| 概念 | 语法/命令 | 用途 | 示例 | 注意事项 |
|---|
| 创建租户 | pulsar-admin tenants create my-tenant | 隔离不同业务或客户 | -a admin 指定管理员 | 支持认证与 ACL |
| 创建命名空间 | pulsar-admin namespaces create my-tenant/my-ns | 租户下逻辑分组 | my-tenant/production | 命名空间是配额管理单位 |
| 设置配额 | pulsar-admin namespaces set-backlog-quota my-tenant/my-ns --limit 10G --policy producer_request_hold | 控制 backlog 大小 | 支持 time、size 限制 | 防止磁盘打满 |
| 设置消息 TTL | pulsar-admin namespaces set-retention my-tenant/my-ns --time 7d --size 100G | 自动清理过期消息 | TTL 7 天或 100GB | 与 backlog 配合使用 |
| 授权管理 | pulsar-admin namespaces grant-permission my-tenant/my-ns --role developer --actions produce,consume | 控制访问权限 | 角色可自定义 | 需配合认证机制 |
- ✅ 多租户优势:资源隔离、独立配额、独立策略、安全可控
第4章:Pulsar 客户端编程(Producer)
4.1 客户端环境搭建(Java / Python / Go)
| 语言 | 依赖/安装命令 | 用途 | 注意事项 |
|---|
| Java | Maven 依赖:org.apache.pulsar:pulsar-client:3.3.0 | 构建 Producer/Consumer | JDK 8+,推荐使用最新稳定版 |
| Python | pip install pulsar-client | Python 客户端 | 需安装 C++ 依赖(libpulsar) |
| Go | go get github.com/apache/pulsar-client-go/pulsar | Go 语言支持 | 支持异步和 TLS |
| Node.js | npm install pulsar-client | JS/TS 支持 | 社区维护,功能较 Java 少 |
- ✅ 所有客户端需连接 brokerServiceUrl(如
pulsar://localhost:6650)
Java Maven 依赖:
<dependency>
<groupId>org.apache.pulsar</groupId>
<artifactId>pulsar-client</artifactId>
<version>3.3.0</version>
</dependency>
4.2 创建 Producer 实例与配置参数
| 配置项 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| serviceUrl | .serviceUrl("pulsar://localhost:6650") | 指定 Broker 地址 | 必填项 | 可为单节点或多节点列表 |
| topicName | .topic("my-topic") | 指定发送的 Topic | Topic 不存在会自动创建 | 建议提前创建 |
| producerName | .producerName("producer-1") | 自定义 Producer 名称 | 便于监控识别 | 非必需 |
| sendTimeout | .sendTimeout(30, TimeUnit.SECONDS) | 发送超时时间 | 防止阻塞 | 网络差时可调大 |
| maxPendingMessages | .maxPendingMessages(1000) | 最大未确认消息数 | 控制内存使用 | 超出将阻塞或丢弃 |
| blockIfQueueFull | .blockIfQueueFull(true) | 队列满时是否阻塞 | false 则抛异常 | 生产环境建议 true |
Java 创建 Producer 示例:
PulsarClient client = PulsarClient.builder()
.serviceUrl("pulsar://localhost:6650")
.build();
Producer<byte[]> producer = client.newProducer()
.topic("my-topic")
.producerName("producer-1")
.sendTimeout(30, TimeUnit.SECONDS)
.maxPendingMessages(1000)
.create();
4.3 发送消息:同步、异步、批量发送
| 发送模式 | 方法 | 用途 | 代码示例 | 注意事项 |
|---|
| 同步发送 | send(msg) | 等待确认后返回 | MessageId msgId = producer.send("Hello".getBytes()); | 影响吞吐,适合关键消息 |
| 异步发送 | sendAsync(msg) | 立即返回 CompletableFuture | 见下方代码块 | 需处理回调,避免内存泄漏 |
| 批量发送 | 启用 batchingEnabled(true) | 多条消息合并发送,提升吞吐 | .enableBatching(true).batchingMaxPublishDelay(10, TimeUnit.MILLISECONDS) | 增加延迟,需权衡吞吐与延迟 |
| 带回调的异步 | sendAsync(msg).whenComplete((msgId, ex) -> {...}) | 处理成功或失败 | 见下方代码块 | 推荐用于生产环境 |
异步发送示例:
producer.sendAsync("Hello".getBytes())
.thenAccept(msgId -> System.out.println("Sent: " + msgId));
带回调的异步发送示例:
producer.sendAsync("Hello".getBytes())
.whenComplete((msgId, ex) -> {
if (ex != null) {
System.err.println("Send failed: " + ex);
} else {
System.out.println("Success: " + msgId);
}
});
4.4 消息路由策略(轮询、按 Key、自定义)
| 路由策略 | 配置方式 | 用途 | 代码示例 | 注意事项 |
|---|
| 轮询(RoundRobin) | 默认策略 | 均匀分布到各分区 | 无需配置 | 适用于无顺序要求场景 |
| 按 Key 路由 | 设置消息 Key | 相同 Key 的消息进入同一分区 | producer.newMessage().key("order-1001").value("data").send(); | 保证 Key 级顺序 |
| 自定义路由 | 实现 MessageRouter 接口 | 完全控制路由逻辑 | 见下方代码块 | 需确保负载均衡 |
| 禁用分区 | 不使用 Partitioned Topic | 所有消息发往单一 Topic | 适用于小流量 | 无并行度 |
自定义路由示例:
.messageRouter((msg, topic, n) -> {
return Math.abs(msg.getKey().hashCode()) % n;
});
4.5 消息压缩与加密(Compression、Encryption)
| 功能 | 配置方法 | 用途 | 代码示例 | 注意事项 |
|---|
| 压缩(Compression) | .compressionType(CompressionType.LZ4) | 减少网络传输大小 | .enableCompression(true).compressionType(CompressionType.ZLIB) | 支持 LZ4、ZLIB、ZSTD、SNAPPY |
| 加密(Encryption) | 配置 CryptoKeyReader | 端到端消息加密 | .addEncryptionKey("key-1").cryptoKeyReader(keyReader) | 需实现密钥读取逻辑 |
| 启用压缩 | enableCompression(true) | 全局启用压缩 | 建议在高吞吐场景开启 | 增加 CPU 开销 |
| 加密算法 | 通过 KeyReader 提供 | 支持 AES、RSA 等 | 需客户端共享密钥 | 安全性高,性能开销大 |
4.6 消息属性与 Schema 支持
| 功能 | 方法 | 用途 | 代码示例 | 注意事项 |
|---|
| 消息属性(Properties) | .property("key", "value") | 添加自定义元数据 | .newMessage().value("data").property("type", "event").send(); | 不超过 1KB,用于路由或过滤 |
| Schema 支持 | .schema(Schema.STRING) | 定义消息数据结构 | Producer<String> producer = client.newProducer(Schema.STRING).topic("my-topic").create(); | 支持 STRING、JSON、AVRO、PROTOBUF |
| JSON Schema | Schema.JSON(MyRecord.class) | 序列化 POJO | 见下方代码块 | 自动生成 Schema 并注册 |
| Schema 兼容性 | 在命名空间设置 | 控制 Schema 演化策略 | pulsar-admin namespaces set-schema-compatibility -c BACKWARD my-tenant/my-ns | 避免消费者解析失败 |
JSON Schema 示例:
public class MyRecord {
public String name;
public int age;
}
// 使用:
Producer<MyRecord> producer = client.newProducer(Schema.JSON(MyRecord.class))
.topic("my-topic")
.create();
- ✅ Schema 优势:类型安全、自动序列化、版本管理、兼容性检查
第5章:Pulsar 客户端编程(Consumer)
5.1 创建 Consumer 实例与订阅模式选择
| 配置项 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| serviceUrl | .serviceUrl("pulsar://localhost:6650") | 指定 Broker 地址 | 必填项,同 Producer | 可为单节点或多节点列表 |
| topic | .topic("my-topic") | 指定消费的 Topic | 支持通配符(*、>) | Topic 必须存在或自动创建 |
| subscriptionName | .subscriptionName("sub-1") | 订阅名称,唯一标识消费者组 | 必填项 | 同一订阅名共享消费位点 |
| subscriptionType | .subscriptionType(SubscriptionType.Exclusive) | 设置订阅模式 | 见 5.2 节详解 | 影响并发和顺序性 |
| consumerName | .consumerName("consumer-1") | 自定义消费者名称 | 便于监控识别 | 非必需 |
| ackTimeout | .ackTimeout(30, TimeUnit.SECONDS) | 消息确认超时时间 | 超时未确认将重投 | 建议设置为处理时间的 2~3 倍 |
| negativeAckRedeliveryDelay | .negativeAckRedeliveryDelay(5, TimeUnit.SECONDS) | Negative Ack 后重试延迟 | 控制错误消息重试频率 | 避免频繁重试压垮系统 |
Java 创建 Consumer 示例:
PulsarClient client = PulsarClient.builder()
.serviceUrl("pulsar://localhost:6650")
.build();
Consumer<byte[]> consumer = client.newConsumer()
.topic("my-topic")
.subscriptionName("sub-1")
.subscriptionType(SubscriptionType.Shared)
.ackTimeout(30, TimeUnit.SECONDS)
.subscribe();
5.2 订阅类型详解:Exclusive、Shared、Failover、Key_Shared
| 订阅类型 | 说明 | 代码设置 | 注意事项 |
|---|
| Exclusive | 一个订阅只能有一个消费者,其他连接将被拒绝 | .subscriptionType(SubscriptionType.Exclusive) | 保证严格顺序,但无高可用 |
| Shared | 多个消费者共享订阅,消息轮询分发(无序) | .subscriptionType(SubscriptionType.Shared) | 适用于任务队列,消费者可动态增减 |
| Failover | 多个消费者注册,按名字排序,主消费者处理,失败后由下一个接管 | .subscriptionType(SubscriptionType.Failover) | 保证顺序性和容错,但并发度为 1 |
| Key_Shared | 消息按 Key 哈希,相同 Key 由同一消费者处理,保证 Key 级顺序 | .subscriptionType(SubscriptionType.Key_Shared) | 需设置 .keySharedPolicy(),支持 sticky 或 auto-split |
Key_Shared 示例:
.subscriptionType(SubscriptionType.Key_Shared)
.keySharedPolicy(KeySharedPolicy.stickyHash())
5.3 消费消息:同步与异步模式
| 消费模式 | 方法 | 用途 | 代码示例 | 注意事项 |
|---|
| 同步消费 | consumer.receive() | 阻塞等待下一条消息 | 见下方代码块 | 简单直观,但吞吐受限 |
| 异步消费 | consumer.receiveAsync() | 返回 CompletableFuture,非阻塞 | 见下方代码块 | 高吞吐,需处理回调 |
| 消息监听器 | messageListener((c, m) -> {...}) | 注册监听器自动接收消息 | 见下方代码块 | 推荐用于生产环境,简化逻辑 |
| 带超时接收 | receive(timeout, unit) | 指定等待时间,避免无限阻塞 | Message<byte[]> msg = consumer.receive(5, TimeUnit.SECONDS); | 适用于定时任务或批处理场景 |
同步消费示例:
Message<byte[]> msg = consumer.receive();
System.out.println("Received: " + new String(msg.getData()));
consumer.acknowledge(msg);
异步消费示例:
consumer.receiveAsync().thenAccept(msg -> {
System.out.println("Received: " + new String(msg.getData()));
consumer.acknowledgeAsync(msg);
});
消息监听器示例:
.messageListener((c, m) -> {
System.out.println("Auto: " + new String(m.getData()));
c.acknowledge(m);
})
5.4 消息确认机制(Individual、Cumulative)
| 确认类型 | 方法 | 用途 | 代码示例 | 注意事项 |
|---|
| Individual Ack | acknowledge(msg) 或 acknowledge(msgId) | 确认单条消息 | consumer.acknowledge(msg); | 精确控制,避免重处理 |
| Cumulative Ack | acknowledgeCumulative(msg) | 确认该消息及之前所有消息 | consumer.acknowledgeCumulative(msg); | 仅支持 Exclusive 和 Failover 模式 |
| 批量确认 | 多条消息逐个或累积确认 | 提升确认效率 | 可结合异步确认 | 避免内存积压 |
| 异步确认 | acknowledgeAsync(msg) | 非阻塞确认,提升吞吐 | 见下方代码块 | 推荐在高并发场景使用 |
异步确认示例:
consumer.acknowledgeAsync(msg).whenComplete((result, ex) -> {
if (ex != null) System.err.println("Ack failed");
});
⚠️ 注意事项:
- Shared 和 Key_Shared 模式仅支持 Individual Ack
- Cumulative Ack 可能导致消息空洞(gap)无法处理
5.5 Negative Acknowledgment 与重试逻辑
| 功能 | 方法 | 用途 | 代码示例 | 注意事项 |
|---|
| Negative Ack | negativeAcknowledge(msg) | 明确告知处理失败,立即重投 | consumer.negativeAcknowledge(msg); | 比超时更快触发重试 |
| 异步 Negative Ack | negativeAcknowledgeAsync(msg) | 非阻塞 Negative Ack | consumer.negativeAcknowledgeAsync(msg); | 提升性能 |
| 重试延迟配置 | negativeAckRedeliveryDelay(5, TimeUnit.SECONDS) | 设置重试间隔 | 在 Consumer 构建时设置 | 避免雪崩 |
| 死信队列(DLQ) | .deadLetterPolicy(DeadLetterPolicy.builder().maxRedeliverCount(5).build()) | 达到最大重试次数后转入 DLQ | 见下方代码块 | 需单独消费 DLQ 处理异常消息 |
| 重试主题(Retry Topic) | 内置机制 | 自动将失败消息发往重试主题,延迟重试 | 默认启用 | 配合 enableRetry 使用 |
开启重试与 DLQ 示例:
Consumer<byte[]> consumer = client.newConsumer()
.topic("my-topic")
.subscriptionName("sub-retry")
.enableRetry(true)
.deadLetterPolicy(DeadLetterPolicy.builder()
.maxRedeliverCount(5)
.deadLetterTopic("my-dlq-topic")
.build())
.subscribe();
5.6 消费者流控(Flow Control)与背压处理
| 参数 | 语法 | 用途 | 说明 | 注意事项 |
|---|
| receiverQueueSize | .receiverQueueSize(1000) | 设置消费者本地接收队列大小 | 控制预取消息数量 | 默认 1000,过大导致内存占用高 |
| 流控自动调整 | Pulsar 自动 | 根据网络和消费速度动态调整 | 基于 TCP 流控 | 无需手动干预 |
| 背压策略 | subscriptionBacklogQuota(Namespace 级) | 控制 backlog 超限时的行为 | 支持 producer_request_hold、drop Newest | 防止消费者跟不上导致 OOM |
| 手动暂停消费 | pause() / resume() | 暂停和恢复消息接收 | consumer.pause(); consumer.resume(); | 适用于批处理或资源紧张场景 |
| 监控消费延迟 | 使用 pulsar-admin topics stats | 查看消费滞后情况 | pulsar-admin topics stats persistent://public/default/my-topic | 及时发现消费瓶颈 |
- ✅ 建议:合理设置
receiverQueueSize,结合 DLQ 和监控,避免背压导致系统崩溃
第6章:Pulsar Topic 与 Namespace 管理
6.1 Topic 类型:Persistent、Non-Persistent、Partitioned
| Topic 类型 | 说明 | 创建方式 | 注意事项 |
|---|
| Persistent | 消息持久化到 BookKeeper,Broker 重启不丢失 | 默认类型,路径以 persistent:// 开头 | 如 persistent://public/default/my-topic |
| Non-Persistent | 消息仅在内存中传输,不写入 BookKeeper | 路径以 non-persistent:// 开头 | 适用于低延迟、可丢失场景(如监控) |
| Partitioned Topic | 一个逻辑 Topic 拆分为多个物理分区,提升吞吐 | 需指定分区数 | 通过 pulsar-admin topics create-partitioned-topic 创建 |
| 分区数限制 | 最大分区数在创建时指定,不可动态增加 | 默认最大 64,可配置 | 建议根据并发需求预估 |
- ✅ 分区 Topic 优势:水平扩展、高吞吐、并行消费
6.2 创建与删除 Topic(命令行与 API)
| 操作 | 命令/API | 用途 | 示例 | 注意事项 |
|---|
| 创建 Persistent Topic | pulsar-admin topics create <topic> | 创建普通持久化 Topic | pulsar-admin topics create persistent://public/default/my-topic | 若命名空间允许,可自动创建 |
| 创建 Partitioned Topic | pulsar-admin topics create-partitioned-topic -p 4 <topic> | 创建 4 分区的 Topic | -p 4 表示 4 个分区 | 分区数不可更改 |
| 删除 Topic | pulsar-admin topics delete <topic> | 删除 Topic 及其数据 | 支持 --force 强制删除 | 删除后数据不可恢复 |
| 列出 Topic | pulsar-admin topics list <namespace> | 查看命名空间下所有 Topic | pulsar-admin topics list public/default | 支持通配符 |
| Java API 创建 | admin.topics().createNonPartitionedTopic(topic) | 编程方式创建 Topic | 需引入 pulsar-admin 依赖 | 适合自动化管理 |
6.3 Namespace 的作用与资源配置
| 功能 | 说明 | 配置命令 | 注意事项 |
|---|
| 资源隔离 | Namespace 是配额、策略、命名的基本单位 | pulsar-admin namespaces create my-tenant/my-ns | 属于某个租户 |
| 配置存储 | 可在 Namespace 级设置 TTL、Retention、Backlog 等 | 统一管理策略 | 避免逐 Topic 配置 |
| 多租户支持 | 一个租户可有多个 Namespace | my-tenant/prod、my-tenant/dev | 实现环境隔离 |
| 复制策略 | 设置跨集群复制 | pulsar-admin namespaces set-clusters --clusters cl1,cl2 my-tenant/my-ns | 需启用 Geo-Replication |
| 角色授权 | 在 Namespace 级授予权限 | pulsar-admin namespaces grant-permission ... | 精细化权限控制 |
6.4 Topic 策略管理(Retention、TTL、Backlog)
| 策略 | 配置命令 | 用途 | 示例 | 注意事项 |
|---|
| Retention 策略 | set-retention | 保留已确认消息的时间或大小 | pulsar-admin namespaces set-retention my-tenant/my-ns --time 7d --size 100G | 用于审计或重放场景 |
| TTL(Time To Live) | set-message-ttl | 未被消费的消息最大存活时间 | pulsar-admin namespaces set-message-ttl my-tenant/my-ns --ttl 3600 | 超时后自动删除 |
| Backlog 配额 | set-backlog-quota | 控制未确认消息的存储上限 | --limit 10G --policy producer_request_hold | 防止磁盘打满 |
| 配额策略 | producer_request_hold、consumer_backlog_eviction、block | 超限时的处理行为 | hold 表示生产者阻塞 | block 会拒绝新消息 |
| 清理 backlog | clear-backlog | 强制清除某个订阅的 backlog | pulsar-admin topics clear-backlog my-topic -s sub1 | 慎用,可能导致消息丢失 |
6.5 分区 Topic 的负载均衡与扩展
| 操作 | 命令/API | 用途 | 说明 | 注意事项 |
|---|
| 查看分区分布 | pulsar-admin topics list-partitioned-topic | 查看分区状态 | 显示每个分区的 Owner Broker | 用于排查不均衡 |
| 负载均衡触发 | pulsar-admin brokers reload 或自动 | 重新分配 Topic 到 Broker | 基于 bundle 负载 | 建议开启动态负载均衡 |
| 扩展分区数 | 不支持动态扩展 | 无法增加已有 Partitioned Topic 的分区数 | 必须重建 Topic | 需提前规划 |
| 分区路由策略 | 轮询、按 Key、自定义 | 决定消息进入哪个分区 | Producer 配置 | 影响并行度和顺序性 |
| 监控分区吞吐 | pulsar-admin topics stats <topic> | 查看每个分区的生产/消费速率 | 识别热点分区 | 可通过 Key 分布优化 |
- ✅ 建议:合理预估分区数,使用 Key_Shared 消费模式提升并行处理能力
第7章:Pulsar 消息存储与持久化
7.1 BookKeeper 架构与 Ledger 机制
| 概念名称 | 说明 | 注意事项 |
|---|
| BookKeeper | 分布式日志存储系统,由多个 Bookie 节点组成,提供低延迟、高可用的日志写入服务 | Pulsar 的持久化层,独立部署 |
| Bookie | BookKeeper 的存储节点,每个 Bookie 管理本地磁盘上的日志文件(Journal、Ledger) | 建议使用 SSD 提升性能 |
| Ensemble | 写入一个 Ledger 时使用的 Bookie 集合,如 ensemble=3 表示从集群中选择 3 个 Bookie | 由客户端(Broker)动态选择 |
| Quorum 机制 | 包括 writeQuorum(写入副本数)和 ackQuorum(确认成功所需副本数),保证数据一致性 | 默认均为 2(在 ensemble=3 时) |
| Ledger Manager | 管理 Ledger 元数据(如位置、状态),元数据存储在 ZooKeeper 中 | Broker 通过 ZooKeeper 查找 Ledger 信息 |
7.2 Entry、Ledger、Bookie 的关系
| 概念 | 说明 | 注意事项 |
|---|
| Entry | 消息在 BookKeeper 中的最小单位,对应一条 Pulsar 消息或批处理中的单条记录 | 包含 Entry ID(递增)、数据体、元数据 |
| Ledger | 逻辑上的日志流,由一系列有序的 Entry 组成,不可变,写满后自动滚动 | 每个 Pulsar Topic 分区对应一个活跃 Ledger |
| Bookie | 物理存储节点,负责存储多个 Ledger 的 Entries | 每个 Bookie 可服务多个 Ledger |
| 映射关系 | 一条消息 → 封装为 Entry → 写入某个 Ledger → 分布存储在多个 Bookie 上 | 一个 Ledger 可跨多个 Bookie 存储 |
🔗 数据路径示例:
Producer → Broker → (Entry) → Ledger-123 → [Bookie-A, Bookie-B, Bookie-C]
7.3 写入流程:从 Producer 到 Bookie
| 步骤 | 说明 | 注意事项 |
|---|
| 1. 生产者发送消息 | Producer 连接 Broker,调用 send() 发送消息 | 支持同步、异步、批量模式 |
| 2. Broker 接收并封装 | Broker 将消息封装为 Entry,确定归属的 Ledger | Ledger 由 Topic 分区决定 |
| 3. 选择 Ensemble | Broker 根据负载和策略选择一组 Bookie(如 3 个)作为 Ensemble | 使用一致性哈希或轮询策略 |
| 4. 并行写入 Bookies | Broker 将 Entry 并行发送到 Ensemble 中的多个 Bookie | 支持批处理提升吞吐 |
| 5. Bookie 持久化 | 每个 Bookie 将 Entry 写入内存缓存和 Journal 文件(事务日志) | Journal 确保崩溃可恢复 |
| 6. 返回 Ack | 当达到 ackQuorum 数量的 Bookie 返回成功时,Broker 向 Producer 确认 | 默认配置下,2/3 成功即返回 |
| 7. 更新 Ledger 元数据 | Broker 更新 Ledger 的最后 Entry ID,元数据写入 ZooKeeper | 确保读取一致性 |
✅ 写入特点:
- 异步持久化,高吞吐
- 分片并行写入,低尾延迟
- Journal + Ledger 文件分离,提升 IO 效率
7.4 数据复制与高可用保障
| 机制 | 说明 | 注意事项 |
|---|
| 多副本写入 | 每条 Entry 同时写入多个 Bookie(由 ensemble 和 writeQuorum 控制) | 默认 2~3 副本,确保数据不丢失 |
| 多数派确认(Quorum) | 只需 ackQuorum 个副本确认即可返回成功,容忍部分节点慢或故障 | 实现高可用与高性能平衡 |
| Ledger 自动恢复 | 若 Bookie 宕机,其他 Bookie 仍可提供服务;重启后自动补全缺失数据 | 通过 ZooKeeper 协调恢复过程 |
| Bookie 故障隔离 | 单个 Bookie 故障不影响其他 Ledger 的读写 | Ensemble 可动态调整 |
| 数据完整性校验 | 每个 Entry 包含 checksum,防止数据损坏 | 在写入和读取时验证 |
| 备份与恢复 | 支持将 Ledger 数据备份到 HDFS 或对象存储 | 结合分层存储实现长期归档 |
- ✅ SLA 保障:即使 1 个 Bookie 故障,仍可正常读写(在 3 节点集群中)
7.5 Ledger 自动滚动与存储优化
| 功能 | 说明 | 注意事项 |
|---|
| Ledger 滚动触发条件 | 达到最大大小(maxEntriesPerLedger)、达到最长时间(ledgerRolloverTimeout)、Broker 重启 | 滚动后创建新 Ledger,旧 Ledger 关闭 |
| 配置参数 | managedLedgerMaxEntriesPerLedger=50000、managedLedgerRolloverTimeInMillis=14400000(4小时) | 可在 broker.conf 中设置 |
| 存储优化:垃圾回收 | 不再引用的 Ledger(无活跃订阅)会被标记为可删除 | 通过 compaction 或 offload 清理 |
| 分层存储(Tiered Storage) | 将冷数据从 BookKeeper 卸载到 S3、GCS 等对象存储 | 降低存储成本 |
| 数据压缩 | 支持 LZ4、ZSTD 等压缩算法减少存储占用 | 在 Bookie 级配置 |
| Ledger 删除策略 | 当所有订阅都消费完某个 Ledger 的数据后,自动删除 | 由 Broker 定期检查 |
- ✅ 运维建议:定期监控 Ledger 数量和大小,合理配置滚动策略
第8章:Pulsar 消息确认与重试机制
8.1 Acknowledgment 原理与模式
| 确认模式 | 说明 | 适用场景 | 注意事项 |
|---|
| Individual Ack | 消费者逐条确认消息,每条消息独立标记为已处理 | Shared、Key_Shared 订阅 | 精确控制,避免重处理 |
| Cumulative Ack | 确认某条消息后,该消息之前的所有消息也被视为已确认 | Exclusive、Failover 订阅 | 提升确认效率,但可能跳过未处理消息 |
| Negative Ack (NACK) | 明确告知消息处理失败,立即触发重投 | 处理异常、依赖不可用 | 比超时更快,避免等待 |
| 异步确认 | acknowledgeAsync() 非阻塞调用,提升吞吐 | 高并发消费场景 | 需处理回调异常 |
| 确认超时(Ack Timeout) | 消息发出后未在指定时间内确认,Broker 自动重投 | 防止消费者卡住 | 设置应大于平均处理时间 |
- ✅ 确认位点(Cursor):消费者组的消费进度保存在 BookKeeper 中,即使消费者重启也能从上次位置继续
8.2 Dead Letter Queue(DLQ)配置与使用
| 配置项 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 启用 DLQ | .deadLetterPolicy(...) | 设置最大重试次数后转入 DLQ | 见下方代码块 | 必须配合 enableRetry(true) |
| 最大重试次数 | maxRedeliverCount=5 | 控制重试上限 | 建议 3~10 次 | 过多可能导致延迟 |
| DLQ Topic 名称 | deadLetterTopic="xxx" | 自定义死信队列名称 | 默认为 {topic}-dlq | 建议集中管理 |
| 消费 DLQ | 使用普通 Consumer 订阅 DLQ Topic | 处理异常消息 | 可人工介入或自动修复后重新投递 | 需监控 DLQ 积压情况 |
DLQ 配置示例:
.deadLetterPolicy(DeadLetterPolicy.builder()
.maxRedeliverCount(5)
.deadLetterTopic("my-dlq-topic")
.build())
✅ 典型流程:正常消费失败 → NACK 或超时 → 重试 → 达到 maxRedeliverCount → 转入 DLQ
8.3 Retry Topic 与延迟重试策略
| 机制 | 说明 | 配置方式 | 注意事项 |
|---|
| Retry Topic | Pulsar 自动生成的重试主题,用于暂存失败消息,支持延迟重投 | 启用 enableRetry(true) 后自动创建 | 格式:{original-topic}-retry-{subscription} |
| 延迟重试策略 | 消息按指数退避或固定间隔重试(如 1s, 5s, 10s, 30s) | 无需额外配置,默认启用 | 避免频繁重试压垮系统 |
| Negative Ack 重试延迟 | negativeAckRedeliveryDelay=5s | 设置 NACK 后的立即重试延迟 | 优先级高于指数退避 |
| 重试消息属性 | Retry Topic 中的消息包含原始主题、重试次数等元数据 | 消费者可通过 getMessage().getProperty(...) 获取 | 用于调试和路由 |
| 关闭重试 | enableRetry(false) | 禁用重试功能,失败后直接进入 DLQ | 适用于不允许重试的业务 |
- ✅ 优势:内置重试机制,无需外部调度器,简化错误处理逻辑
8.4 消费失败处理与容错设计
| 场景 | 处理策略 | 推荐做法 | 注意事项 |
|---|
| 瞬时异常(网络抖动、DB 连接超时) | 自动重试(Retry Topic) | 启用 enableRetry(true) | 设置合理的重试间隔 |
| 永久性错误(数据格式错误、逻辑异常) | 转入 DLQ | 配置 maxRedeliverCount | 避免无限重试 |
| 消费者宕机 | Broker 检测到连接断开,自动重分配 | 使用 Shared 或 Key_Shared 模式 | 确保高可用 |
| 消息积压(Backlog) | 增加消费者实例、优化处理逻辑 | 监控 backlog 大小 | 防止 OOM |
| 顺序性要求 | 使用 Key_Shared 订阅 + Individual Ack | 确保同 Key 消息由同一消费者处理 | 避免并发导致乱序 |
| 幂等消费 | 业务层保证重复处理不产生副作用 | 使用数据库唯一键、Redis 记录已处理 ID | 因重试和 NACK 可能导致重复 |
✅ 容错设计原则:
- 快速失败:异常时及时 NACK
- 分级重试:指数退避 + DLQ
- 可观测性:监控重试、DLQ、消费延迟
- 人工干预通道:提供 DLQ 消费和重放能力
第9章:Pulsar Schema 与数据格式
9.1 Schema 的作用与类型(String、JSON、Avro、Protobuf)
| 类型 | 说明 | 注意事项 |
|---|
| String | 纯文本格式,适用于简单字符串消息 | 编码默认 UTF-8 |
| JSON | 使用 JSON 格式存储对象,可自动序列化/反序列化 POJO | 支持嵌套结构,可读性强 |
| Avro | 二进制格式,Schema 定义在 JSON 中,高效压缩和序列化 | 需预定义 Schema |
| Protobuf | Google 开源的高效二进制序列化格式,需 .proto 文件定义结构 | 性能高,体积小 |
| Schema 的作用 | 类型安全、自动序列化/反序列化、兼容性检查、消费者无需关心数据格式 | 减少出错,提升开发效率 |
- ✅ Schema 优势:避免”消息格式地狱”,实现生产者与消费者解耦
9.2 Schema 注册与版本管理
| 操作 | 命令/API | 用途 | 示例 | 注意事项 |
|---|
| 自动注册 | Producer 首次发送时自动注册 | 无需手动干预 | 启用 isEnableSchemaValidation=true | 建议在开发环境使用 |
| 手动注册 | pulsar-admin schemas upload | 提前上传 Schema 定义 | pulsar-admin schemas upload --filename user.avsc my-topic | 适合生产环境控制变更 |
| 查看 Schema | pulsar-admin schemas get <topic> | 获取当前 Schema 信息 | 返回 JSON 格式的 Schema 定义 | 用于调试和审计 |
| 删除 Schema | pulsar-admin schemas delete <topic> | 清除 Topic 的 Schema | 删除后可重新注册 | 谨慎操作,影响现有消息 |
| Schema 版本 | 系统自动为每次变更生成版本号(如 v1, v2) | 跟踪 Schema 演化历史 | 存储在 BookKeeper 中 | 不可回滚,只能向前演进 |
✅ 注册流程:Producer 发送消息 → Broker 检查 Schema 是否存在 → 不存在则注册 → 存在则校验兼容性
9.3 Schema 兼容性策略(BACKWARD、FORWARD、FULL)
| 兼容性策略 | 说明 | 适用场景 | 注意事项 |
|---|
| BACKWARD | 新 Schema 可读旧数据(推荐) | 消费者升级先于生产者 | 最常用,保证向下兼容 |
| FORWARD | 旧 Schema 可读新数据(新增字段可为空) | 生产者升级先于消费者 | 新增字段必须为 optional |
| FULL | 新旧 Schema 双向兼容 | 严格要求双向兼容 | 限制最多,适用于核心系统 |
| NONE | 不检查兼容性,允许任意变更 | 临时测试或完全重构 | 风险高,可能导致消费失败 |
配置方式:在 Namespace 级设置:
pulsar-admin namespaces set-schema-compatibility -c BACKWARD public/default
默认兼容性为 BACKWARD。
✅ 兼容性示例(JSON/Avro):
- ✅ 允许:添加 optional 字段、修改字段默认值
- ❌ 禁止:删除字段、修改字段类型(如
string → int)
9.4 使用 Schema 定义消息结构(Producer 与 Consumer)
Java 示例:使用 JSON Schema
| 配置项 | 语法 | 用途 | 代码示例 | 注意事项 |
|---|
| 定义 POJO | public class User { public String name; public int age; } | 消息数据结构 | 必须有默认构造函数 | 字段需 public 或提供 getter/setter |
| Producer 使用 Schema | .schema(Schema.JSON(User.class)) | 自动序列化对象 | 见下方代码块 | Schema 自动注册 |
| Consumer 使用 Schema | .schema(Schema.JSON(User.class)) | 自动反序列化为对象 | 见下方代码块 | 类定义必须一致或兼容 |
Producer 使用 Schema:
Producer<User> producer = client.newProducer(Schema.JSON(User.class))
.topic("user-topic")
.create();
producer.send(new User("Alice", 30));
Consumer 使用 Schema:
Consumer<User> consumer = client.newConsumer(Schema.JSON(User.class))
.topic("user-topic")
.subscriptionName("sub1")
.subscribe();
Message<User> msg = consumer.receive();
User user = msg.getValue(); // 直接获取对象
Avro Schema 示例
// 定义 Avro Schema
String schemaJson = "{"
+ " \"type\": \"record\","
+ " \"name\": \"User\","
+ " \"fields\": ["
+ " {\"name\": \"name\", \"type\": \"string\"},"
+ " {\"name\": \"age\", \"type\": \"int\"}"
+ " ]"
+ "}";
Schema<User> schema = Schema.AVRO(User.class, schemaJson);
- ✅ 跨语言兼容:Java 生产 JSON Schema 消息,Python 消费者可直接使用
pulsar.schema.JSONSchema(User) 反序列化
第10章:Pulsar Functions(轻量级计算)
10.1 Pulsar Functions 概述与使用场景
| 概念 | 说明 | 注意事项 |
|---|
| Pulsar Function | 轻量级无服务器(Serverless)计算单元,处理 Pulsar 消息流 | 无需管理服务器 |
| 输入 Topic | Function 消费的消息源 | 可为单个或多个 Topic |
| 输出 Topic | Function 处理后发送结果的目标 Topic | 可指定或动态路由 |
| 使用场景 | 实时数据清洗、协议转换、聚合统计、路由分发、调用外部服务 | 替代小型 Flink/Spark Job |
| 资源隔离 | 每个 Function 独立运行,失败不影响其他 | 支持命名空间级资源限制 |
- ✅ 核心优势:简单、轻量、与 Pulsar 深度集成
10.2 开发第一个 Function(Java/Python/Go)
Java Function 示例
| 步骤 | 说明 | 代码示例 | 注意事项 |
|---|
| 实现 Function 接口 | 继承 java.util.function.Function 或实现 PulsarFunction | 见下方代码块 | Context 提供日志、状态等能力 |
| 打包 | 使用 Maven 构建 Fat JAR | 包含所有依赖 | mvn clean package |
| 部署 | 使用 pulsar-admin functions 命令 | 见 10.3 节 | 类名需完整 |
public class ExclamationFunction implements Function<String, String> {
@Override
public String process(String input, Context context) {
return input + "!";
}
}
Python Function 示例
def process(input_data):
return input_data.upper()
pulsar-admin functions create \
--py myfunc.py \
--classname process \
--tenant public \
--namespace default \
--name upper-func \
--inputs persistent://public/default/in \
--output persistent://public/default/out
Go Function 示例
package main
import "context"
func HandleRequest(ctx context.Context, in []byte) ([]byte, error) {
return []byte(string(in) + "!"), nil
}
使用 go build 编译为二进制部署。
10.3 部署模式:LocalRun 与 Cluster 模式
| 部署模式 | 说明 | 命令示例 | 注意事项 |
|---|
| LocalRun 模式 | 在本地运行 Function,用于测试和调试 | 见下方代码块 | 不注册到集群,重启后消失 |
| Cluster 模式 | 将 Function 提交到 Pulsar 集群,由 Worker 节点调度执行 | 见下方代码块 | 支持高可用、监控、自动恢复 |
| 更新 Function | update 命令替换代码 | pulsar-admin functions update --name myfunc --jar new-version.jar | 保持名称一致 |
| 删除 Function | delete 命令移除 | pulsar-admin functions delete --name myfunc | 停止运行并清除元数据 |
LocalRun 模式:
pulsar-admin functions localrun \
--jar target/my-function.jar \
--className com.example.MyFunction \
--inputs my-input-topic \
--output my-output-topic
Cluster 模式:
pulsar-admin functions create \
--jar target/my-function.jar \
--name myfunc \
--inputs in-topic \
--output out-topic
- ✅ Cluster 模式依赖:需启用
functionsWorkerEnabled=true 并配置 Worker 节点
10.4 状态存储(State API)与有状态处理
| 功能 | 说明 | 代码示例(Java) | 注意事项 |
|---|
| 状态存储 | 每个 Function 可维护键值状态,跨消息持久化 | context.putState("count", count + 1); context.getState("key"); | 状态存储在内部 Topic 中 |
| 状态一致性 | 提供恰好一次(exactly-once)语义保证 | 需启用 processingGuarantees=ATLEAST_ONCE 或 EFFECTIVELY_ONCE | 后者性能较低 |
| 状态用途 | 计数器、缓存、聚合窗口、去重 | 适合有状态流处理 | 避免存储过大对象 |
| 清理状态 | 手动删除或设置 TTL | context.deleteState("key"); | 防止状态无限增长 |
| 状态后端 | 默认使用 Pulsar 内部 Topic 存储 | 无需外部数据库 | 高可用,自动复制 |
计数器 Function 示例:
public String process(String input, Context context) {
Integer count = context.getState("count");
count = (count == null) ? 0 : count + 1;
context.putState("count", count);
return input + " [" + count + "]";
}
10.5 并行度与容错机制
| 特性 | 说明 | 配置方式 | 注意事项 |
|---|
| 并行度(Parallelism) | 将 Function 拆分为多个实例并行处理消息 | --parallelism 3 | 提升吞吐,需消息可并行处理 |
| 实例(Instance) | 每个并行单元称为 Instance,独立运行 | 由 Worker 动态调度 | 故障时自动迁移 |
| 容错机制 | Instance 失败后由 Worker 自动重启或迁移 | 基于 ZooKeeper 协调 | 保证至少一次处理 |
| 消息重试 | 处理失败时自动重试(可配置) | 结合 Negative Ack | 避免无限重试 |
| 处理保证 | 支持 ATLEAST_ONCE、EFFECTIVELY_ONCE、ATMOST_ONCE | --processing-guarantees EFFECTIVELY_ONCE | 后者需启用状态检查点 |
| 资源隔离 | 每个 Instance 分配独立资源(CPU、内存) | 在 functionsWorker 配置 | 防止资源争抢 |
✅ 建议:
- 无状态 Function:高并行度 + ATLEAST_ONCE
- 有状态 Function:合理并行度 + EFFECTIVELY_ONCE
第11章:Pulsar IO(Connectors)
11.1 Pulsar IO 架构与工作原理
| 概念名称 | 说明 | 注意事项 |
|---|
| Pulsar IO | Pulsar 的数据集成框架,用于在 Pulsar 和外部系统之间移动数据 | 类似 Kafka Connect |
| Worker 节点 | 运行 Connectors 的进程,可独立部署或嵌入 Broker | 由 functionsWorker 提供支持 |
| Source Connector | 从外部系统读取数据并写入 Pulsar Topic | 如 MySQL、Kafka、文件 |
| Sink Connector | 从 Pulsar Topic 读取消息并写入外部系统 | 如 Elasticsearch、HBase、JDBC |
| Connector 实例 | 每个 Connector 可运行多个并行实例 | 提升吞吐 |
| 配置存储 | Connector 配置保存在 ZooKeeper 中 | 支持动态更新 |
✅ 工作流程:外部系统 → Source Connector → Pulsar Topic → Sink Connector → 目标系统
11.2 Source 连接器(Kafka、MySQL、File)
| Connector 类型 | 用途 | 配置示例(YAML) | 注意事项 |
|---|
| Kafka Source | 从 Kafka 集群消费数据并写入 Pulsar | 见下方代码块 | 需 Kafka 客户端权限 |
| MySQL Source | 监听 MySQL Binlog,捕获数据变更(CDC) | 见下方代码块 | 需启用 Binlog,使用 row 格式 |
| File Source | 读取本地或远程文件(如日志文件)并发送到 Pulsar | 见下方代码块 | 支持轮询或 inotify 监听 |
Kafka Source:
tenant: public
namespace: default
name: kafka-source
topicName: pulsar-topic
builtin: kafka-source
configs:
bootstrapServers: "kafka-broker:9092"
groupId: pulsar-group
topic: kafka-topic
MySQL Source:
builtin: mysql-source
configs:
hostname: mysql-host
port: 3306
database: mydb
username: user
password: pass
tableName: users
File Source:
builtin: file-source
configs:
directory: "/logs"
filePattern: ".*\\.log"
includeHidden: false
部署命令:
pulsar-admin sources create \
--archive connectors/pulsar-io-kafka-source.nar \
--name kafka-source \
--destination-topic-name my-topic ...
11.3 Sink 连接器(Elasticsearch、HBase、JDBC)
| Connector 类型 | 用途 | 配置示例(YAML) | 注意事项 |
|---|
| Elasticsearch Sink | 将消息写入 ES,用于搜索和分析 | 见下方代码块 | 支持批量写入,提升性能 |
| HBase Sink | 写入 HBase 表,适用于海量结构化存储 | 见下方代码块 | 需配置 HBase 客户端资源 |
| JDBC Sink | 将消息写入关系型数据库(MySQL、PostgreSQL 等) | 见下方代码块 | 支持字段映射,需 Schema 匹配 |
Elasticsearch Sink:
tenant: public
namespace: default
name: es-sink
inputTopics: persistent://public/default/logs
builtin: elasticsearch-sink
configs:
elasticSearchUrl: "http://es:9200"
indexName: log-index
typeName: log
batchSize: 100
HBase Sink:
builtin: hbase-sink
configs:
hbaseConfigResources: "hbase-site.xml"
tableName: "user_table"
rowKeyColumn: "id"
columnsMap: '{"name": "cf:name", "age": "cf:age"}'
JDBC Sink:
builtin: jdbc-sink
configs:
userName: "user"
password: "pass"
jdbcUrl: "jdbc:mysql://db:3306/mydb"
tableName: "events"
- ✅ 通用配置项:
inputTopics、processingGuarantees、parallelism、maxPendingAsyncRequests
11.4 自定义 Connector 开发流程
| 步骤 | 说明 | 代码示例(Java) | 注意事项 |
|---|
| 1. 继承 Source/Sink 接口 | 实现 Source<T> 或 Sink<T> 接口 | 见下方代码块 | read() 方法返回 Record<T> |
| 2. 实现核心方法 | open() 初始化,read() 读取数据,close() 释放资源 | 见下方代码块 | 支持同步或异步写入 |
| 3. 打包为 NAR | 使用 Maven 构建 NAR(Nested ARchive)包 | 见下方代码块 | NAR 包含所有依赖 |
| 4. 注册与部署 | 将 NAR 文件上传并创建 Connector | pulsar-admin sources create --archive my-connector.nar ... | 类名需正确 |
| 5. 配置文件 | 提供 config.yaml 定义可配置参数 | 便于用户自定义 | |
实现 Source 接口:
public class CustomSource implements Source<String> {
@Override
public void open(Map<String, Object> config, SourceContext context) throws Exception { ... }
@Override
public Record<String> read() throws Exception { ... }
}
实现 Sink 接口:
public class CustomSink implements Sink<String> {
@Override
public void write(Record<String> record) throws Exception {
// 写入外部系统
}
}
Maven NAR 打包配置:
<packaging>nar</packaging>
<plugin>
<groupId>org.apache.nifi</groupId>
<artifactId>nifi-nar-maven-plugin</artifactId>
</plugin>
配置文件示例:
configs:
apiEndpoint: "string"
batchSize: "int"
- ✅ 建议:参考官方 Connector 示例(如 Kafka、File)进行开发
11.5 Connector 配置与监控
| 功能 | 语法/命令 | 用途 | 示例 | 注意事项 |
|---|
| 创建 Connector | pulsar-admin sources/sinks create | 部署 Source 或 Sink | --name my-sink --inputs my-topic --archive connector.nar | 支持 --brokerServiceUrl |
| 更新 Connector | update 命令 | 修改配置或升级版本 | pulsar-admin sinks update --name my-sink --parallelism 3 | 自动滚动重启 |
| 删除 Connector | delete 命令 | 移除并停止运行 | pulsar-admin sources delete --name my-source | 数据流将中断 |
| 查看状态 | get-status 命令 | 检查运行状态和实例健康度 | pulsar-admin sinks get-status --name my-sink | 显示运行实例数、失败次数 |
| 监控指标 | Prometheus + Grafana | 采集吞吐、延迟、错误率 | 指标前缀:pulsar_io_* | 需启用 exposeConsumerLevelMetrics |
| 日志查看 | pulsar-admin functions logs | 查看 Connector 实例日志 | --name my-connector --instance-id 0 | 日志存储在 Worker 节点 |
- ✅ 运维建议:定期检查状态、监控 backlog、设置告警规则
第12章:Pulsar 与生态集成
12.1 与 Flink 集成:流处理管道
| 功能 | 说明 | 代码示例(Java) | 注意事项 |
|---|
| Flink Pulsar Source | 作为 Flink 数据源消费 Pulsar 消息 | 见下方代码块 | 需引入 pulsar-flink-connector |
| Flink Pulsar Sink | 将 Flink 处理结果写回 Pulsar | 见下方代码块 | 支持精确一次语义 |
| 一致性保证 | 支持 EXACTLY_ONCE 和 AT_LEAST_ONCE | 在 PulsarSourceOptions 中配置 | 需启用 Checkpoint |
| 优势 | 低延迟、高吞吐、无缝集成 | 替代 Kafka + Flink 架构 | 充分利用 Pulsar 分层存储 |
Flink Pulsar Source:
StreamExecutionEnvironment env = StreamExecutionEnvironment.getExecutionEnvironment();
PulsarSource<String> source = PulsarSource.builder()
.serviceUrl("pulsar://localhost:6650")
.adminUrl("http://localhost:8080")
.topic("my-topic")
.subscriptionName("flink-sub")
.deserializationSchema(PulsarDeserializationSchema.flinkSchema(Schema.STRING))
.build();
DataStream<String> stream = env.fromSource(source,
WatermarkStrategy.noWatermarks(), "Pulsar Source");
Flink Pulsar Sink:
PulsarSink<String> sink = PulsarSink.builder()
.topic("output-topic")
.producerConfigurator(ProducerConfigurationData.builder().build())
.serializationSchema(PulsarSerializationSchema.flinkSchema(Schema.STRING))
.serviceUrl("pulsar://localhost:6650")
.adminUrl("http://localhost:8080")
.build();
stream.sinkTo(sink);
- ✅ 适用场景:实时 ETL、复杂事件处理(CEP)、实时数仓
12.2 与 Spark Streaming 集成
| 集成方式 | 说明 | 代码示例(Scala) | 注意事项 |
|---|
| 使用 Pulsar Spark Connector | 官方或社区提供的 Spark DataSource | 见下方代码块 | 需添加依赖 io.streamnative.connectors:spark-streaming-pulsar_2.12 |
| Streaming 模式 | 支持微批处理(Micro-batching) | 基于 Spark Structured Streaming | 延迟高于 Flink |
| Schema 支持 | 自动映射 Pulsar Schema 到 Spark DataFrame | 支持 JSON、Avro | 需类型匹配 |
| 性能调优 | 配置批大小、并行度、缓存策略 | --conf spark.streaming.pulsar.receiver.maxRatePerPartition=1000 | 避免反压 |
Spark Streaming 示例:
val df = spark.readStream
.format("pulsar")
.option("service.url", "pulsar://localhost:6650")
.option("admin.url", "http://localhost:8080")
.option("topic", "my-topic")
.load()
df.select($"value" cast "string").writeStream
.outputMode("append")
.format("console")
.start()
.awaitTermination()
- ✅ 适用场景:批流一体、机器学习管道、历史数据分析
12.3 与 Kafka 兼容层(Pulsar-Kafka API)
| 功能 | 说明 | 配置/代码示例 | 注意事项 |
|---|
| Kafka on Pulsar(KoP) | 在 Pulsar Broker 上启用 Kafka 协议支持 | 见下方配置 | 无需额外网关 |
| 生产者兼容 | 使用 Kafka Producer 写入 Pulsar Topic | 见下方代码块 | Topic 自动映射到 Pulsar |
| 消费者兼容 | 使用 Kafka Consumer 消费 Pulsar 消息 | 见下方代码块 | 支持 Offset 管理 |
| 限制 | 不支持 Kafka Streams、部分高级特性缺失 | 如事务、Exactly-Once 跨 Topic | 适用于迁移过渡期 |
broker.conf 配置:
enableKafkaBroker=true
kafkaAdvertisedListeners=PLAINTEXT://localhost:9092
Kafka 生产者兼容:
Properties props = new Properties();
props.put("bootstrap.servers", "localhost:9092");
props.put("key.serializer", "org.apache.kafka.common.serialization.StringSerializer");
props.put("value.serializer", "org.apache.kafka.common.serialization.StringSerializer");
KafkaProducer<String, String> producer = new KafkaProducer<>(props);
producer.send(new ProducerRecord<>("kop-topic", "Hello Pulsar"));
Kafka 消费者兼容:
props.put("bootstrap.servers", "localhost:9092");
props.put("group.id", "kop-group");
KafkaConsumer<String, String> consumer = new KafkaConsumer<>(props);
consumer.subscribe(Arrays.asList("kop-topic"));
- ✅ 优势:零代码迁移 Kafka 应用,逐步切换到原生 Pulsar API
12.4 与 Spring Boot 集成开发
| 集成方式 | 说明 | 配置示例(YAML) | 注意事项 |
|---|
| 使用 Spring Pulsar Starter | 官方或社区 Starter 简化集成 | 见下方代码块 | 依赖 spring-pulsar-starter |
| 发送消息(Producer) | 使用 PulsarTemplate | 见下方代码块 | 支持泛型和 Schema |
| 接收消息(Consumer) | 使用 @PulsarListener 注解 | 见下方代码块 | 支持同步/异步消费 |
| Schema 支持 | 自动处理 POJO 序列化 | 配合 Schema.JSON(User.class) | 需定义实体类 |
| 自动配置 | Spring Boot 自动创建 Client、Producer、Consumer | 基于配置属性 | 支持 Profile 切换环境 |
application.yml 配置:
pulsar:
client:
service-url: pulsar://localhost:6650
admin-url: http://localhost:8080
producer:
default-producer-name: spring-producer
consumer:
default-subscription-name: spring-sub
发送消息:
@Autowired
private PulsarTemplate<String> template;
template.send("my-topic", "Hello Spring");
接收消息:
@PulsarListener(subscriptionName = "sub1", topics = "my-topic")
public void listen(String message) {
System.out.println("Received: " + message);
}
- ✅ 开发优势:快速构建微服务,与 Spring Cloud 生态无缝集成
第13章:Pulsar 安全机制
13.1 认证机制(TLS、JWT、OAuth2)
| 认证方式 | 说明 | 配置/代码示例 | 注意事项 |
|---|
| TLS 证书认证 | 使用 X.509 证书验证客户端身份 | 见下方配置 | 需 CA 签发,适合内部系统 |
| JWT 认证 | 基于 JSON Web Token 的无状态认证 | 见下方配置 | 需管理密钥,适合微服务 |
| OAuth2 认证 | 集成外部身份提供商(如 Keycloak、Google) | 见下方配置 | 支持 SSO,适合多租户平台 |
| 启用认证 | Broker 端开启认证链 | authenticationEnabled=true,多种方式可共存 | 需在 broker.conf 中设置 |
| 客户端配置 | 指定认证方式连接集群 | 见下方代码块 | Token 可从文件或 URL 获取 |
broker.conf 配置(TLS 证书认证):
brokerClientTlsEnabled=true
tlsCertificateFilePath=/path/to/broker.crt
tlsKeyFilePath=/path/to/broker.key
客户端配置 tlsTrustCertsFilePath。
JWT 认证:
# 生成 Token
pulsar tokens create --secret-key file://secret.key --subject user1
broker.conf 配置(OAuth2):
authenticationProviders=org.apache.pulsar.broker.authentication.AuthenticationProviderOAuth2
oauth2IssuerUrl=https://auth.example.com/realms/pulsar
oauth2Audience=pulsar-cluster
认证链配置(broker.conf):
authenticationEnabled=true
authenticationProviders=org.apache.pulsar.broker.authentication.AuthenticationProviderJwt,org.apache.pulsar.broker.authentication.AuthenticationProviderTls
客户端连接示例:
client = PulsarClient.builder()
.serviceUrl("pulsar+ssl://broker:6651")
.authentication("token", "bearer-token-string")
.build();
- ✅ 推荐组合:生产环境使用 TLS + JWT,确保传输与身份双重安全
13.2 授权与 ACL 管理
| 功能 | 配置命令 | 用途 | 示例 | 注意事项 |
|---|
| 启用授权 | authorizationEnabled=true | 开启 ACL 控制 | 在 broker.conf 中设置 | 默认关闭 |
| 授予命名空间权限 | pulsar-admin namespaces grant-permission | 控制租户下命名空间访问 | --actions produce,consume --role user1 public/default | 支持 produce、consume、functions 等 |
| 授予 Topic 权限 | pulsar-admin topics grant-permission | 细粒度控制单个 Topic | persistent://public/default/my-topic --role user2 --actions produce | 优先级高于命名空间 |
| 查看权限 | pulsar-admin namespaces permissions <namespace> | 审计当前 ACL 规则 | 显示角色与操作映射 | 用于安全合规检查 |
| 移除权限 | revoke-permission 命令 | 撤销用户访问权限 | pulsar-admin namespaces revoke-permission public/default --role guest | 立即生效 |
| 默认策略 | superUserRoles | 定义超级用户(绕过 ACL) | superUserRoles=admin,monitor | 建议最小化配置 |
- ✅ 权限模型:Tenant → Namespace → Topic 三级继承,子级可覆盖父级策略
13.3 加密传输与静态数据加密
| 加密类型 | 说明 | 配置方式 | 注意事项 |
|---|
| 传输加密(TLS) | 所有客户端与 Broker 之间通信加密 | brokerServicePortTls=6651、webServicePortTls=8443、启用 tlsEnabledInBroker=true | 需配置证书和信任链 |
| 内部组件加密 | Broker ↔ Bookie、Broker ↔ Broker 通信加密 | bookieClientTlsEnabled=true、clusterTlsEnabled=true | 防止内网嗅探 |
| 静态数据加密 | Bookie 磁盘数据加密(需插件支持) | 使用 EncryptedLedgerStorage | 增加 CPU 开销,需密钥管理 |
| 客户端配置 | 客户端必须信任服务器证书 | .tlsTrustCertsFilePath("/path/to/ca.crt").enableTls(true) | 否则连接失败 |
| 通配符证书 | 支持多域名或 IP | 确保证书中包含所有 Broker 地址 | 避免 hostname 不匹配错误 |
- ✅ 合规要求:金融、医疗等敏感行业必须启用全链路加密
13.4 多租户隔离与资源配额
| 特性 | 说明 | 配置示例 | 注意事项 |
|---|
| 多租户模型 | 支持多个租户(Tenant)独立管理命名空间 | pulsar-admin tenants create tenant1 --admin-roles user1 | 租户间完全隔离 |
| 命名空间配额 | 限制每个命名空间的资源使用 | 见下方代码块 | 防止资源耗尽 |
| 生产/消费速率限制 | 控制 Producer 和 Consumer 的吞吐 | --dispatch-rate-period(秒)、--rate-period(限流周期) | 单位:msg/s 或 byte/s |
| 存储配额 | 限制命名空间最大存储容量 | set-storage-quota --quota 100G public/default | 超额后 Producer 被阻塞 |
| 隔离策略 | 为租户分配专用 Broker | 见下方代码块 | 实现物理隔离 |
| 监控与告警 | 配合 Prometheus 监控配额使用情况 | 指标:pulsar_backlog_size、pulsar_storage_used | 设置阈值告警 |
命名空间配额:
pulsar-admin namespaces set-backlog-quota public/default \
--limit 10G --policy producer_request_hold
pulsar-admin namespaces set-subscription-dispatch-rate public/default \
--msg-dispatch-rate 1000
隔离策略:
pulsar-admin namespaces set-isolation-policy public/default \
--auto-failover-policy-type min_available \
--auto-failover-policy-params "min_limit=1,broker_count=3"
- ✅ 最佳实践:为每个业务线创建独立租户,避免”邻居干扰”
第14章:性能调优与最佳实践
14.1 生产环境配置建议
| 组件 | 推荐配置 | 说明 | 注意事项 |
|---|
| Broker | brokerServicePort=6650、webServicePort=8080、numExecutorThreads=16 | 根据 CPU 核心数调整线程池 | 避免过度竞争 |
| Bookie | useV2WireProtocol=true、journalSyncData=true、flushInterval=100(ms) | 提升写入性能和可靠性 | Journal 建议使用 SSD |
| ZooKeeper | 独立部署 3 或 5 节点集群 | 不与 Broker/Bookie 共用 | 保证协调服务高可用 |
| JVM | -Xms8g -Xmx8g、-XX:+UseG1GC、-XX:MaxGCPauseMillis=10 | 减少 GC 停顿 | Bookie 建议更大堆内存 |
| 磁盘规划 | Journal 和 Ledger 目录分离 | Journal 使用 SSD,Ledger 可使用 HDD | 避免 IO 争抢 |
| 网络 | 10Gbps 网络,低延迟交换机 | 满足高吞吐需求 | 避免跨机房写入 |
- ✅ 部署架构:Broker 无状态可横向扩展,Bookie 存储层独立部署
14.2 Producer 与 Consumer 性能参数调优
Producer 调优
| 参数 | 配置项 | 推荐值 | 说明 | 注意事项 |
|---|
| 批处理大小 | batchingMaxMessages=1000 | 根据消息大小调整 | 提升吞吐,降低请求开销 | 过大会增加延迟 |
| 批处理超时 | batchingMaxPublishDelay=1ms | 平衡吞吐与延迟 | 达到时间或大小即发送 | 延迟敏感业务设为 0 |
| 压缩算法 | compressionType=LZ4 | 可选 ZSTD, SNAPPY | 减少网络传输量 | CPU 开销增加 |
| 异步发送 | sendAsync() | 推荐使用 | 非阻塞,高吞吐 | 需处理回调异常 |
Consumer 调优
| 参数 | 配置项 | 推荐值 | 说明 | 注意事项 |
|---|
| 接收队列大小 | receiverQueueSize=1000 | 控制预取数量 | 避免 OOM | 高吞吐场景可调大 |
| 消费模式 | subscriptionType=Shared 或 Key_Shared | 根据业务选择 | Exclusive 无并发 | — |
| 确认超时 | ackTimeoutMillis=30000 | 防止消息卡住 | 应大于平均处理时间 | — |
| 负载均衡 | consumerEventListener | 监听分区变化 | 用于动态调整 | — |
- ✅ 建议:使用异步 API + 合理批处理,避免同步阻塞
14.3 BookKeeper 性能优化(Write Cache、Flush Interval)
| 优化项 | 配置参数 | 推荐值 | 说明 | 注意事项 |
|---|
| 写缓存大小 | managedLedgerCacheSize=2g | 建议 1/4~1/3 物理内存 | 缓存热数据,提升读性能 | 过大会影响 GC |
| Journal 同步 | journalSyncData=true | 确保数据不丢失 | 写入磁盘前不返回 | 可设 false 提升性能(有风险) |
| Ledger Flush 间隔 | flushInterval=100(ms) | 控制内存数据落盘频率 | 减少随机写 | 过长可能丢失最近数据 |
| Direct Memory | directMemoryLimit=4g | Bookie 使用堆外内存 | 减少 GC 压力 | 需监控 off-heap 使用 |
| V2 协议 | useV2WireProtocol=true | 更高效的数据传输协议 | 推荐启用 | 需所有节点一致 |
| 磁盘监控 | diskUsageThreshold=0.8 | 设置磁盘使用率阈值 | 超过后停止写入 | 避免磁盘写满 |
- ✅ 性能测试建议:使用 workload-producer 工具压测,观察
pulsar_bookie_* 指标
14.4 背压处理与流量控制
| 机制 | 说明 | 配置方式 | 注意事项 |
|---|
| Producer 端背压 | 当 Broker 缓冲区满时阻塞或失败 | blockIfQueueFull=true 或 false | 设为 false 可快速失败 |
| Consumer 端流控 | Broker 根据接收队列大小控制发送速率 | 自动启用,无需配置 | 基于 receiverQueueSize |
| Backlog 配额 | 限制未确认消息总量 | set-backlog-quota --limit 10G --policy producer_request_hold | 策略可选 producer_request_hold 或 consumer_backlog_eviction |
| 分层存储卸载 | 将冷数据卸载到对象存储 | trigger=threshold 或手动触发 | 释放 Bookie 存储压力 |
| 监控指标 | pulsar_backlog_size、pulsar_in_bytes_total | 使用 Prometheus 采集 | 设置告警规则 |
| 自动扩缩容 | 结合 Kubernetes HPA | 基于 CPU 或 backlog 扩容 Consumer | 实现弹性伸缩 |
✅ 处理策略:
- 短时积压:增加 Consumer 实例
- 长期积压:优化处理逻辑或启用分层存储
14.5 大消息处理与批处理策略
| 场景 | 策略 | 配置/代码示例 | 注意事项 |
|---|
| 大消息(>5MB) | 启用 Chunking(分块) | 见下方代码块 | Consumer 需启用 autoAckOldestChunkedMessageOnQueueFull |
| 批处理 | 合并多条消息为 Batch | .batchingMaxMessages(1000).batchingMaxPublishDelay(10, TimeUnit.MILLISECONDS) | 推荐大小 1~2MB |
| 压缩 | 减小消息体积 | .compressionType(Compression.ZSTD) | 适合文本类数据 |
| 分层存储 | 大消息自动归档 | 配置 managedLedgerOffloadThreshold | 降低实时存储成本 |
| 消息大小限制 | 控制单条消息上限 | maxMessageSize=5242880(5MB) | 防止 OOM |
| 替代方案 | 大文件存对象存储,消息仅传 URL | {"fileName": "report.pdf", "url": "s3://bucket/report.pdf"} | 推荐用于超大文件 |
Chunking 示例:
producer = client.newProducer()
.topic("large-topic")
.enableChunking(true)
.maxPendingChunks(10)
.create();
✅ 最佳实践:
- 普通消息:≤ 1MB,启用批处理 + 压缩
- 大消息:> 1MB,使用 Chunking 或外存引用
第15章:高级特性与扩展
15.1 延迟消息(Delayed Delivery)
| 概念名称 | 说明 | 注意事项 |
|---|
| 延迟消息 | 消息发送后不立即投递给消费者,而是在指定延迟时间后才可被消费 | 适用于定时任务、重试机制 |
| deliverAfter | 指定从发送时刻起延迟一定时间 | 单位:毫秒、秒等 |
| deliverAt | 指定消息在某个绝对时间点后投递 | 如 System.currentTimeMillis() + 60_000 |
| 时间轮算法 | Pulsar 使用时间轮(Timing Wheel)高效管理延迟消息 | 内存中维护延迟队列 |
| 仅共享/Key-Shared 订阅支持 | Exclusive 和 Failover 模式不支持延迟消息 | 因独占消费无法等待 |
| 配置项 | 语法/配置 | 用途 | 代码示例(Java) | 注意事项 |
|---|
| 设置延迟时间 | deliverAfter(long delay, TimeUnit unit) | 发送延迟消息 | 见下方代码块 | 延迟时间从 send() 调用开始计算 |
| 绝对时间投递 | deliverAt(long timestamp) | 在指定时间点后投递 | .deliverAt(System.currentTimeMillis() + 30000) | 时间戳为毫秒级 Unix 时间 |
| Broker 控制 | maxMessageDelay=30d | 限制最大延迟时间 | 在 broker.conf 中设置 | 默认 1 天,避免无限堆积 |
| 消费者行为 | 自动过滤未到期消息 | 消费者不可见延迟中消息 | 无需额外配置 | — |
设置延迟时间:
producer.newMessage()
.value("Delayed Message")
.deliverAfter(10, TimeUnit.SECONDS)
.send();
- ✅ 典型场景:订单超时取消、消息重试退避、定时通知
15.2 消息去重(Message Deduplication)
| 概念名称 | 说明 | 注意事项 |
|---|
| 消息去重 | 防止 Producer 重复发送相同消息导致消费者重复处理 | 保证”恰好一次”语义 |
| Producer ID | 每个 Producer 实例的唯一标识 | 用于跟踪消息序列 |
| Sequence ID | 每条消息的递增序号 | 与 Producer ID 组合判断重复 |
| 元数据存储 | 去重信息存储在 BookKeeper 中 | 包含 producer_id、sequence_id、event_time |
| 时间窗口 | 仅在指定时间范围内检查重复 | 默认 6 小时 |
| 配置项 | 语法/配置 | 用途 | 配置示例 | 注意事项 |
|---|
| 启用去重 | brokerDeduplicationEnabled=true | Broker 级开启去重功能 | 在 broker.conf 中设置 | 默认关闭 |
| 命名空间级启用 | set-deduplication | 为命名空间开启去重 | pulsar-admin namespaces set-deduplication public/default --enable | 更灵活控制 |
| Producer 去重 | enableProducerDeduplication(true) | Producer 声明支持去重 | .enableProducerDeduplication(true) | 必须设置,否则无效 |
| 去重时间窗口 | deduplicationMaxTime=1h | 设置重复检查的时间范围 | 单位:s、m、h、d | 过长影响性能 |
| 去重状态清理 | 自动清理过期元数据 | 基于 deduplicationMaxTime | 无需手动干预 | — |
✅ 去重流程:Producer 发送 → Broker 检查 (Producer ID + Sequence ID) → 已存在则丢弃 → 否则写入
15.3 分层存储(Tiered Storage)
| 概念名称 | 说明 | 注意事项 |
|---|
| 分层存储 | 将冷数据从 BookKeeper 卸载到低成本对象存储(如 S3、GCS、HDFS) | 无限存储消息,降低成本 |
| 卸载策略 | 基于时间或大小触发卸载 | 如 offloadAfterThreshold=100G |
| 支持的存储后端 | AWS S3、Google GCS、Aliyun OSS、HDFS、Azblob | 需配置对应连接器 |
| 数据访问透明性 | 消费者可无缝读取已卸载消息 | Broker 自动从远端拉取 |
| 元数据一致性 | Ledger 元数据仍保留在 ZooKeeper | 确保数据可定位 |
| 配置项 | 语法/配置 | 用途 | 配置示例 | 注意事项 |
|---|
| 启用分层存储 | managedLedgerOffloadDeletionLagMs=10m | 设置卸载延迟(防止未确认消息被卸载) | 在 broker.conf 中设置 | 必须大于最大 ack 超时 |
| 存储类型 | tieredStoragePrimaryType=S3 | 指定远端存储类型 | 可选 S3, GCS, HDFS 等 | — |
| S3 配置 | 多个 S3 相关参数 | 连接 AWS S3 | 见下方代码块 | 建议使用 IAM 角色 |
| 触发卸载 | pulsar-admin topics offload | 手动卸载 | --size-threshold 1G my-topic | 查看状态使用 offload-status |
| 自动卸载 | managedLedgerOffloadThreshold=1073741824 | 达到大小后自动卸载 | 单位:字节 | 需配合 offloadDeletionLag |
S3 配置(broker.conf):
awsAccessKey=xxx
awsSecretKey=yyy
s3ManagedLedgerOffloadRegion=us-west-2
s3ManagedLedgerOffloadBucket=my-pulsar-bucket
- ✅ 优势:实现”热数据高速访问 + 冷数据低成本存储”的混合架构
15.4 事件时间与水位线支持
| 概念名称 | 说明 | 注意事项 |
|---|
| 事件时间(Event Time) | 消息在生产系统中生成的时间,而非到达 Pulsar 的时间 | 用于精确流处理 |
| 水位线(Watermark) | 表示事件时间的进度,用于处理乱序事件 | 控制窗口计算的触发时机 |
| 乱序容忍 | 允许一定时间窗口内的乱序消息 | 如 allowedLateness=5m |
| 与 Flink 集成 | Pulsar 提供事件时间,Flink 用于窗口聚合 | 实现精确的实时统计 |
| 时间语义 | 支持 EventTime、IngestionTime、ProcessingTime | Pulsar 主要提供前两者 |
| 配置项 | 语法/配置 | 用途 | 代码示例(Java) | 注意事项 |
|---|
| 设置事件时间 | eventTime(long timestamp) | Producer 设置事件时间 | 见下方代码块 | 时间戳为毫秒 |
| 获取事件时间 | message.getEventTime() | Consumer 读取事件时间 | long eventTime = msg.getEventTime(); | 用于业务逻辑或转发 |
| Schema 支持 | 使用 Avro/JSON Schema 包含时间字段 | 结构化时间信息 | 如 {"timestamp": 1678886400000, "data": "..."} | 更易处理 |
| 水位线生成 | 通常由 Flink 或 Pulsar Functions 生成 | 基于事件时间分布 | 需自定义逻辑 | Pulsar 本身不直接管理水位线 |
设置事件时间:
producer.newMessage()
.value("data")
.eventTime(System.currentTimeMillis())
.send();
- ✅ 适用场景:实时风控、会话窗口统计、基于时间的 ETL
15.5 Pulsar Proxy 与边缘接入
| 概念名称 | 说明 | 注意事项 |
|---|
| Pulsar Proxy | 位于客户端与 Broker 之间的代理服务,统一接入入口 | 隐藏内部 Broker 拓扑 |
| TLS 终止 | Proxy 可集中处理 TLS 加密解密 | 减轻 Broker 负担 |
| 认证代理 | Proxy 验证客户端身份后转发请求 | 支持 JWT、TLS、OAuth2 |
| 多租户接入 | 为不同租户提供独立接入点 | 结合 DNS 或负载均衡 |
| 边缘部署 | 在边缘节点部署 Proxy,就近接入 | 降低跨区域延迟 |
| 连接复用 | Proxy 复用与 Broker 的连接 | 减少 Broker 连接数压力 |
| 配置项 | 语法/配置 | 用途 | 配置示例 | 注意事项 |
|---|
| 启动 Proxy | bin/pulsar proxy | 启动 Proxy 服务 | 独立脚本 | 需配置 proxy.conf |
| Broker 连接 | brokerServiceUrl=pulsar://broker:6650 | Proxy 连接后端 Broker | 可配置多个 | 支持 DNS 轮询 |
| Web 服务端口 | webServicePort=8080 | HTTP 接口端口 | 默认 8080 | 用于 Admin API |
| 认证启用 | authenticationEnabled=true | 开启客户端认证 | 需配置 authenticationProviders | 与 Broker 一致 |
| TLS 配置 | tlsEnabled=true | 启用 HTTPS 和加密 | 配置证书路径 | 客户端连接使用 pulsar+ssl:// |
✅ 部署模式:Client → Pulsar Proxy(公网)→ 内网 Broker 集群,实现安全隔离与弹性扩展
附录 A:pulsar-admin 常用命令速查表
集群管理
| 命令语法 | 用途说明 | 示例 | 注意事项 |
|---|
clusters list | 列出所有集群 | pulsar-admin clusters list | 多集群部署时使用 |
clusters get <cluster> | 查看集群配置 | pulsar-admin clusters get standalone | 显示 brokerUrl、serviceUrl 等 |
clusters create <name> --broker-url | 创建新集群 | pulsar-admin clusters create my-cluster --broker-url pulsar://broker:6650 | 需提前规划拓扑 |
租户管理
| 命令语法 | 用途说明 | 示例 | 注意事项 |
|---|
tenants create <tenant> --admin-roles | 创建租户并指定管理员角色 | pulsar-admin tenants create myteam --admin-roles admin-user | 租户是资源隔离基础 |
tenants list | 列出所有租户 | pulsar-admin tenants list | 用于审计和管理 |
tenants delete <tenant> | 删除租户 | pulsar-admin tenants delete guest | 必须先删除所有命名空间 |
命名空间管理
| 命令语法 | 用途说明 | 示例 | 注意事项 |
|---|
namespaces create <tenant/namespace> | 创建命名空间 | pulsar-admin namespaces create public/default | 命名空间是配置单位 |
namespaces list <tenant> | 列出某租户下的命名空间 | pulsar-admin namespaces list public | — |
namespaces delete <ns> | 删除命名空间 | pulsar-admin namespaces delete public/temp | 必须无活跃 Topic |
Topic 管理
| 命令语法 | 用途说明 | 示例 | 注意事项 |
|---|
topics list <namespace> | 列出命名空间下所有 Topic | pulsar-admin topics list public/default | 支持持久化/非持久化 |
topics create <topic> -p <partitions> | 创建分区 Topic | pulsar-admin topics create persistent://public/default/news -p 6 | 分区数不可变 |
topics delete <topic> | 删除 Topic | pulsar-admin topics delete persistent://public/default/old-topic | 数据将永久丢失 |
topics stats <topic> | 查看 Topic 实时统计 | pulsar-admin topics stats persistent://public/default/news | 包含吞吐、延迟、backlog |
topics peek-messages <topic> -s <sub> | 查看订阅中第一条消息 | pulsar-admin topics peek-messages -t persistent://public/default/q1 -s sub1 -n 1 | 用于故障排查 |
权限管理
| 命令语法 | 用途说明 | 示例 | 注意事项 |
|---|
namespaces grant-permission <ns> --role --actions | 授予命名空间权限 | pulsar-admin namespaces grant-permission public/default --role user1 --actions produce,consume | 权限:produce, consume, functions |
topics grant-permission <topic> --role --actions | 授予 Topic 级权限 | pulsar-admin topics grant-permission persistent://public/default/secure-topic --role analyst --actions consume | 细粒度控制 |
namespaces permissions <ns> | 查看命名空间权限 | pulsar-admin namespaces permissions public/default | 审计使用 |
资源配额
| 命令语法 | 用途说明 | 示例 | 注意事项 |
|---|
namespaces set-backlog-quota <ns> --limit --policy | 设置 backlog 配额 | pulsar-admin namespaces set-backlog-quota public/default --limit 10G --policy producer_request_hold | 防止磁盘耗尽 |
namespaces set-subscription-dispatch-rate <ns> --msg-dispatch-rate | 限制消费速率 | pulsar-admin namespaces set-subscription-dispatch-rate public/default --msg-dispatch-rate 500 | 单位:msg/s |
namespaces set-publish-rate <ns> --msg-publish-rate | 限制生产速率 | pulsar-admin namespaces set-publish-rate public/default --msg-publish-rate 1000 | 防止突发流量冲击 |
Broker 状态
| 命令语法 | 用途说明 | 示例 | 注意事项 |
|---|
brokers list <cluster> | 列出集群中所有 Broker | pulsar-admin brokers list standalone | 查看活跃节点 |
broker-stats monitoring-metrics | 导出 Prometheus 监控指标 | pulsar-admin broker-stats monitoring-metrics | 可重定向到文件 |
broker-stats load-report | 查看 Broker 负载报告 | pulsar-admin broker-stats load-report | 包含 CPU、内存、带宽 |
Functions 管理
| 命令语法 | 用途说明 | 示例 | 注意事项 |
|---|
functions create --archive --tenant --namespace --name | 创建 Pulsar Function | pulsar-admin functions create --archive my-func.nar --name my-func --inputs persistent://public/default/in --output persistent://public/default/out | 需 NAR 包 |
functions list | 列出所有函数 | pulsar-admin functions list | 查看运行状态 |
functions logs --name --instance-id | 查看函数日志 | pulsar-admin functions logs --name my-func --instance-id 0 | 用于调试 |
Sources/Sinks
| 命令语法 | 用途说明 | 示例 | 注意事项 |
|---|
sources create --archive --name --destination-topic | 创建 Source 连接器 | pulsar-admin sources create --archive connectors/pulsar-io-kafka-source.nar --name kafka-in --destination-topic pulsar-topic --broker-service-url pulsar://localhost:6650 | 需指定 broker URL |
sinks create --archive --name --inputs | 创建 Sink 连接器 | pulsar-admin sinks create --archive connectors/pulsar-io-es-sink.nar --name es-out --inputs my-topic --broker-service-url pulsar://localhost:6650 | — |
sinks get-status --name | 查看 Sink 运行状态 | pulsar-admin sinks get-status --name es-out | 检查是否 healthy |
- ✅ 提示:
- 所有命令可通过
--help 查看详细参数
- 建议将
pulsar-admin 封装为脚本用于自动化运维
- 生产环境操作前建议先在测试环境验证
附录 B:Pulsar 监控指标(Prometheus)详解
Broker 通用
| 指标名称 | 用途说明 | 示例查询(PromQL) | 注意事项 |
|---|
pulsar_broker_count | 当前活跃 Broker 数量 | pulsar_broker_count{cluster="my-cluster"} | 用于集群健康检查 |
pulsar_jvm_memory_bytes_used | JVM 内存使用量 | rate(pulsar_jvm_memory_bytes_used{area="heap"}[5m]) | 监控堆内存增长 |
pulsar_jvm_gc_seconds_count | GC 次数 | rate(pulsar_jvm_gc_seconds_count[1m]) > 10 | 频繁 GC 需优化 JVM |
Topic 级吞吐
| 指标名称 | 用途说明 | 示例查询(PromQL) | 注意事项 |
|---|
pulsar_topics_count | 命名空间下 Topic 数量 | pulsar_topics_count{namespace="public/default"} | 防止 Topic 泛滥 |
pulsar_msg_in_total | 消息进入总数(生产) | rate(pulsar_msg_in_total[1m]) | 单位:条/秒 |
pulsar_msg_out_total | 消息发出总数(消费) | rate(pulsar_msg_out_total[1m]) | 对比 msg_in 判断积压 |
pulsar_throughput_in_bytes_total | 入流量字节数 | rate(pulsar_throughput_in_bytes_total[1m]) / 1024 / 1024 | 单位:MB/s |
pulsar_throughput_out_bytes_total | 出流量字节数 | rate(pulsar_throughput_out_bytes_total[1m]) | 网络带宽监控 |
延迟与积压
| 指标名称 | 用途说明 | 示例查询(PromQL) | 注意事项 |
|---|
pulsar_entry_serialize_latency_le_... | 消息序列化延迟分布 | histogram_quantile(0.99, sum(rate(pulsar_entry_serialize_latency_ms_bucket[5m])) by (le)) | P99 延迟告警 |
pulsar_entry_publish_latency_le_... | 消息发布延迟(从 send 到持久化) | histogram_quantile(0.95, sum(rate(pulsar_entry_publish_latency_ms_bucket[5m])) by (le)) | 核心性能指标 |
pulsar_backlog_size | 未确认消息积压量(字节) | pulsar_backlog_size{topic="persistent://public/default/news"} | 持续增长需扩容 Consumer |
pulsar_subscription_unacked_messages | 未确认消息数 | pulsar_subscription_unacked_messages{subscription="sub1"} | 判断消费滞后 |
存储与 BookKeeper
| 指标名称 | 用途说明 | 示例查询(PromQL) | 注意事项 |
|---|
pulsar_storage_write_latency_le_... | 存储写入延迟 | histogram_quantile(0.99, sum(rate(pulsar_storage_write_latency_le_...[5m])) by (le)) | 反映 Bookie 性能 |
pulsar_storage_size | 命名空间存储使用量 | pulsar_storage_size{namespace="public/default"} | 用于配额告警 |
pulsar_ledger_write_qps | Ledger 写入 QPS | rate(pulsar_ledger_write_qps[1m]) | Bookie 负载指标 |
pulsar_bookie_journal_write_data_size | Journal 写入大小 | rate(pulsar_bookie_journal_write_data_size[1m]) | SSD 写入压力 |
连接与请求
| 指标名称 | 用途说明 | 示例查询(PromQL) | 注意事项 |
|---|
pulsar_connection_count | 当前客户端连接数 | pulsar_connection_count{broker="broker-1"} | 判断连接泄漏 |
pulsar_rate_in | Broker 入请求速率 | rate(pulsar_rate_in[1m]) | 包括 produce/consume 请求 |
pulsar_rate_out | Broker 出请求速率 | rate(pulsar_rate_out[1m]) | — |
pulsar_pending_bytes | 待处理数据量 | pulsar_pending_bytes{namespace="public/default"} | 反映处理压力 |
Functions 与 Connectors
| 指标名称 | 用途说明 | 示例查询(PromQL) | 注意事项 |
|---|
pulsar_function_received_total | Function 接收消息数 | rate(pulsar_function_received_total[1m]) | 判断函数处理能力 |
pulsar_function_processed_successfully_total | 成功处理消息数 | rate(pulsar_function_processed_successfully_total[1m]) | 计算失败率 |
pulsar_sink_written_bytes_total | Sink 写出字节数 | rate(pulsar_sink_written_bytes_total[1m]) | 验证数据导出正常 |
pulsar_source_num_read_from_source | Source 读取外部系统次数 | rate(pulsar_source_num_read_from_source[1m]) | 检查 Source 活跃度 |
分层存储
| 指标名称 | 用途说明 | 示例查询(PromQL) | 注意事项 |
|---|
pulsar_offloaded_bytes_total | 已卸载到远端存储的字节数 | pulsar_offloaded_bytes_total{topic="..."} | 验证分层存储生效 |
pulsar_offload_error_total | 卸载错误次数 | rate(pulsar_offload_error_total[1m]) > 0 | 配置错误或网络问题 |
pulsar_managed_ledger_offload_threshold | 触发卸载的阈值 | pulsar_managed_ledger_offload_threshold | 用于对比实际使用量 |
✅ 监控实践建议:
- 使用 Grafana 导入官方 Pulsar Dashboard(ID: 10003)
- 设置关键告警:
pulsar_backlog_size > 10G
pulsar_entry_publish_latency_ms > 100ms(P99)
rate(pulsar_sink_write_failures[5m]) > 0
- 定期审查指标,优化资源配置