Article

消息系统 Pulsar

更新于:2026-07-14

第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 PulsarApache 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 → BookKeeperBroker 将消息封装为 Entry,写入 BookKeeper 的 Ledger。写入是异步的,支持批处理提升吞吐。
BookKeeper 存储多个 Bookie 组成 Ensemble,消息被分片并多副本存储。默认 ackQuorum=2,writeQuorum=2,ensemble=3。
消费者 ← Broker消费者从 Broker 拉取消息,Broker 从 BookKeeper 读取数据并推送。读取路径:Broker → Bookie → Consumer。
元数据 ← ZooKeeperBroker 启动时从 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:分布式日志存储原理

概念说明注意事项
BookieBookKeeper 的存储节点,负责存储 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 等标准路径
启动 Standalonebin/pulsar standalone默认启用 Broker、BookKeeper、ZooKeeper、Proxy
访问管理界面浏览器打开 http://localhost:8080Web UI 端口为 8080,Broker 服务端口 6650
停止服务Ctrl+Cbin/pulsar-daemon stop standalone推荐使用 daemon 模式后台运行
  • ✅ 适用场景:本地开发、测试、学习
  • ❌ 不适用:生产环境(无高可用、性能受限)

3.2 集群模式部署(Docker / Kubernetes / Bare Metal)

部署方式步骤概要注意事项
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 核心配置

配置项语法用途示例值注意事项
brokerServiceUrlbrokerServiceUrl=pulsar://localhost:6650Broker 服务地址(客户端连接)pulsar://broker1:6650多节点需配置为域名或 VIP
webServiceUrlwebServiceUrl=http://localhost:8080HTTP 接口地址(Admin API)http://broker1:8080用于 REST 管理接口
zookeeperServerszookeeperServers=localhost:2181连接本地 ZooKeeper 地址zk1:2181,zk2:2181必须可达
configurationStoreServersconfigurationStoreServers=localhost:2181配置存储 ZooKeeper 地址同上多集群时指向全局 ZooKeeper
managedLedgerDefaultEnsembleSizemanagedLedgerDefaultEnsembleSize=3默认写入的 Bookie 数量3必须 ≤ Bookie 总数
managedLedgerDefaultWriteQuorummanagedLedgerDefaultWriteQuorum=2写入成功需确认的副本数2≤ ensembleSize
managedLedgerDefaultAckQuorummanagedLedgerDefaultAckQuorum=2返回成功前需确认的副本数2≤ writeQuorum

bookkeeper.conf 核心配置

配置项语法用途示例值注意事项
bookiePortbookiePort=3181Bookie 服务端口3181需防火墙开放
zkServerszkServers=localhost:2181连接 ZooKeeper同上必须与 Pulsar 共享
journalDirectoryjournalDirectory=/pulsar/journal事务日志存储路径自定义路径建议使用独立高速磁盘
ledgerDirectoriesledgerDirectories=/pulsar/ledgers数据文件存储路径多路径可用逗号分隔提升 IO 并行度
useV2WireProtocoluseV2WireProtocol=true启用 V2 协议提升性能true建议开启

zk.conf 配置要点

配置项说明注意事项
dataDirZooKeeper 数据存储目录需定期清理 snapshot
clientPort客户端连接端口(默认 2181)多实例部署时需区分
tickTime心跳间隔(毫秒)默认 2000,不建议修改
initLimitFollower 初始化同步时限通常 10
syncLimitFollower 与 Leader 同步时限通常 5

3.4 监控与指标收集(Prometheus + Grafana)

组件指标类型用途配置方式注意事项
Pulsar BrokerJVM、Topic 数、生产/消费速率、延迟监控服务健康与负载broker.conf 中启用:prometheusStatsHttpPort=8080exposeTopicLevelMetricsInPrometheus=true需配置 Prometheus scrape_configs
BookKeeperBookie 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.loglogs/pulsar-broker-*.logBroker 运行日志启动失败、Topic 创建异常检查 ZooKeeper 连接、端口占用
bookkeeper.loglogs/pulsar-bookkeeper-*.logBookie 写入日志写入超时、Ledger 错误检查磁盘空间、journal 同步性能
zookeeper.loglogs/pulsar-zookeeper-*.logZK 选举与事务日志Leader 选举失败、Session 超时检查网络延迟、tickTime 配置
GC 日志logs/gc.logJVM 垃圾回收情况频繁 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 限制防止磁盘打满
设置消息 TTLpulsar-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)

语言依赖/安装命令用途注意事项
JavaMaven 依赖:org.apache.pulsar:pulsar-client:3.3.0构建 Producer/ConsumerJDK 8+,推荐使用最新稳定版
Pythonpip install pulsar-clientPython 客户端需安装 C++ 依赖(libpulsar)
Gogo get github.com/apache/pulsar-client-go/pulsarGo 语言支持支持异步和 TLS
Node.jsnpm install pulsar-clientJS/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")指定发送的 TopicTopic 不存在会自动创建建议提前创建
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 SchemaSchema.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 Ackacknowledge(msg)acknowledge(msgId)确认单条消息consumer.acknowledge(msg);精确控制,避免重处理
Cumulative AckacknowledgeCumulative(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 AcknegativeAcknowledge(msg)明确告知处理失败,立即重投consumer.negativeAcknowledge(msg);比超时更快触发重试
异步 Negative AcknegativeAcknowledgeAsync(msg)非阻塞 Negative Ackconsumer.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_holddrop 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 Topicpulsar-admin topics create <topic>创建普通持久化 Topicpulsar-admin topics create persistent://public/default/my-topic若命名空间允许,可自动创建
创建 Partitioned Topicpulsar-admin topics create-partitioned-topic -p 4 <topic>创建 4 分区的 Topic-p 4 表示 4 个分区分区数不可更改
删除 Topicpulsar-admin topics delete <topic>删除 Topic 及其数据支持 --force 强制删除删除后数据不可恢复
列出 Topicpulsar-admin topics list <namespace>查看命名空间下所有 Topicpulsar-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 配置
多租户支持一个租户可有多个 Namespacemy-tenant/prodmy-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_holdconsumer_backlog_evictionblock超限时的处理行为hold 表示生产者阻塞block 会拒绝新消息
清理 backlogclear-backlog强制清除某个订阅的 backlogpulsar-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 的持久化层,独立部署
BookieBookKeeper 的存储节点,每个 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,确定归属的 LedgerLedger 由 Topic 分区决定
3. 选择 EnsembleBroker 根据负载和策略选择一组 Bookie(如 3 个)作为 Ensemble使用一致性哈希或轮询策略
4. 并行写入 BookiesBroker 将 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=50000managedLedgerRolloverTimeInMillis=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 TopicPulsar 自动生成的重试主题,用于暂存失败消息,支持延迟重投启用 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
ProtobufGoogle 开源的高效二进制序列化格式,需 .proto 文件定义结构性能高,体积小
Schema 的作用类型安全、自动序列化/反序列化、兼容性检查、消费者无需关心数据格式减少出错,提升开发效率
  • ✅ Schema 优势:避免”消息格式地狱”,实现生产者与消费者解耦

9.2 Schema 注册与版本管理

操作命令/API用途示例注意事项
自动注册Producer 首次发送时自动注册无需手动干预启用 isEnableSchemaValidation=true建议在开发环境使用
手动注册pulsar-admin schemas upload提前上传 Schema 定义pulsar-admin schemas upload --filename user.avsc my-topic适合生产环境控制变更
查看 Schemapulsar-admin schemas get <topic>获取当前 Schema 信息返回 JSON 格式的 Schema 定义用于调试和审计
删除 Schemapulsar-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

配置项语法用途代码示例注意事项
定义 POJOpublic 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 消息流无需管理服务器
输入 TopicFunction 消费的消息源可为单个或多个 Topic
输出 TopicFunction 处理后发送结果的目标 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 节点调度执行见下方代码块支持高可用、监控、自动恢复
更新 Functionupdate 命令替换代码pulsar-admin functions update --name myfunc --jar new-version.jar保持名称一致
删除 Functiondelete 命令移除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_ONCEEFFECTIVELY_ONCE后者性能较低
状态用途计数器、缓存、聚合窗口、去重适合有状态流处理避免存储过大对象
清理状态手动删除或设置 TTLcontext.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 IOPulsar 的数据集成框架,用于在 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"
  • ✅ 通用配置项:inputTopicsprocessingGuaranteesparallelismmaxPendingAsyncRequests

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 文件上传并创建 Connectorpulsar-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 配置与监控

功能语法/命令用途示例注意事项
创建 Connectorpulsar-admin sources/sinks create部署 Source 或 Sink--name my-sink --inputs my-topic --archive connector.nar支持 --brokerServiceUrl
更新 Connectorupdate 命令修改配置或升级版本pulsar-admin sinks update --name my-sink --parallelism 3自动滚动重启
删除 Connectordelete 命令移除并停止运行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 与生态集成

功能说明代码示例(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细粒度控制单个 Topicpersistent://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=6651webServicePortTls=8443、启用 tlsEnabledInBroker=true需配置证书和信任链
内部组件加密Broker ↔ Bookie、Broker ↔ Broker 通信加密bookieClientTlsEnabled=trueclusterTlsEnabled=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_sizepulsar_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 生产环境配置建议

组件推荐配置说明注意事项
BrokerbrokerServicePort=6650webServicePort=8080numExecutorThreads=16根据 CPU 核心数调整线程池避免过度竞争
BookieuseV2WireProtocol=truejournalSyncData=trueflushInterval=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=SharedKey_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 MemorydirectMemoryLimit=4gBookie 使用堆外内存减少 GC 压力需监控 off-heap 使用
V2 协议useV2WireProtocol=true更高效的数据传输协议推荐启用需所有节点一致
磁盘监控diskUsageThreshold=0.8设置磁盘使用率阈值超过后停止写入避免磁盘写满
  • ✅ 性能测试建议:使用 workload-producer 工具压测,观察 pulsar_bookie_* 指标

14.4 背压处理与流量控制

机制说明配置方式注意事项
Producer 端背压当 Broker 缓冲区满时阻塞或失败blockIfQueueFull=truefalse设为 false 可快速失败
Consumer 端流控Broker 根据接收队列大小控制发送速率自动启用,无需配置基于 receiverQueueSize
Backlog 配额限制未确认消息总量set-backlog-quota --limit 10G --policy producer_request_hold策略可选 producer_request_holdconsumer_backlog_eviction
分层存储卸载将冷数据卸载到对象存储trigger=threshold 或手动触发释放 Bookie 存储压力
监控指标pulsar_backlog_sizepulsar_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=trueBroker 级开启去重功能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、ProcessingTimePulsar 主要提供前两者
配置项语法/配置用途代码示例(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 连接数压力
配置项语法/配置用途配置示例注意事项
启动 Proxybin/pulsar proxy启动 Proxy 服务独立脚本需配置 proxy.conf
Broker 连接brokerServiceUrl=pulsar://broker:6650Proxy 连接后端 Broker可配置多个支持 DNS 轮询
Web 服务端口webServicePort=8080HTTP 接口端口默认 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>列出命名空间下所有 Topicpulsar-admin topics list public/default支持持久化/非持久化
topics create <topic> -p <partitions>创建分区 Topicpulsar-admin topics create persistent://public/default/news -p 6分区数不可变
topics delete <topic>删除 Topicpulsar-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>列出集群中所有 Brokerpulsar-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 Functionpulsar-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_usedJVM 内存使用量rate(pulsar_jvm_memory_bytes_used{area="heap"}[5m])监控堆内存增长
pulsar_jvm_gc_seconds_countGC 次数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_qpsLedger 写入 QPSrate(pulsar_ledger_write_qps[1m])Bookie 负载指标
pulsar_bookie_journal_write_data_sizeJournal 写入大小rate(pulsar_bookie_journal_write_data_size[1m])SSD 写入压力

连接与请求

指标名称用途说明示例查询(PromQL)注意事项
pulsar_connection_count当前客户端连接数pulsar_connection_count{broker="broker-1"}判断连接泄漏
pulsar_rate_inBroker 入请求速率rate(pulsar_rate_in[1m])包括 produce/consume 请求
pulsar_rate_outBroker 出请求速率rate(pulsar_rate_out[1m])
pulsar_pending_bytes待处理数据量pulsar_pending_bytes{namespace="public/default"}反映处理压力

Functions 与 Connectors

指标名称用途说明示例查询(PromQL)注意事项
pulsar_function_received_totalFunction 接收消息数rate(pulsar_function_received_total[1m])判断函数处理能力
pulsar_function_processed_successfully_total成功处理消息数rate(pulsar_function_processed_successfully_total[1m])计算失败率
pulsar_sink_written_bytes_totalSink 写出字节数rate(pulsar_sink_written_bytes_total[1m])验证数据导出正常
pulsar_source_num_read_from_sourceSource 读取外部系统次数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
  • 定期审查指标,优化资源配置