一、引言:为什么需要 Apache Pulsar?
在消息中间件领域,Apache Kafka 已经成为事实上的标准,但在面对多租户、跨区域复制、分层存储等场景时,其架构局限性也日益暴露。Apache Pulsar 作为一个云原生分布式消息流平台,从一开始就采用了计算与存储分离的架构,天然适配容器化部署和弹性扩缩容场景。本文将从 Pulsar 的核心架构原理出发,深入剖析其分层设计、消息传递模型、Topic 体系与多租户隔离机制,并覆盖生产级部署、性能调优与运维监控的完整工程实践。
二、Pulsar 核心架构:计算与存储分离
Pulsar 的架构可以简洁地分为三层:
Broker 层(计算层):无状态的服务节点,负责协议接入、消息路由、负载均衡、Schema 校验、Geo-Replication 等逻辑。Broker 不存储持久化数据,故障时可在秒级恢复。
BookKeeper 层(存储层):基于预写日志(WAL)的分布式存储引擎。Pulsar 将所有 Topic 的持久化数据下沉到 BookKeeper 的 Bookie 节点。BookKeeper 通过多副本复制、Quorum 写入、Minor Compaction 等机制保证数据可靠性与持久性。
ZooKeeper 层(元数据存储):负责集群协调、Topic 元数据、Leader 选举、配置管理等。Pulsar 2.10+ 版本引入了 Metadata Store 抽象,支持 etcd 等替代方案。
# 查看 Pulsar 集群状态
pulsar-admin brokers list useast
# 查看 Namespace 下所有 Topic
pulsar-admin topics list useast/tenant/ns1
# 查看 Bookie 节点状态
pulsar-admin bookies list-bookies
这种分层架构使得 Pulsar 具备以下核心优势:弹性扩缩容(Broker 可任意扩缩,不影响数据)、存储计算独立扩展、多租户隔离(Namespace 级别)、分层卸载(Tiered Storage 自动卸载冷数据到对象存储)。
三、Namespace 与 Topic 体系
Pulsar 的 Topic 命名空间采用 persistent://tenant/namespace/topic 的 URI 结构,这种层级设计是其多租户能力的基石。
Namespace 是 Pulsar 中策略管理的基本单位,所有策略——如消息保留策略、TTL、Backlog Quota、Auth、Geo-Replication、Dispatch Rate 等——都在 Namespace 级别统一配置。这意味着同一个 Namespace 下的所有 Topic 共享同一组治理策略。
# 创建 Namespace 并设置保留策略
pulsar-admin namespaces create useast/tenant/prod
pulsar-admin namespaces set-retention useast/tenant/prod \
--size 10G \
--time 7d
# 设置 Dispatch Rate 限流
pulsar-admin namespaces set-dispatch-rate useast/tenant/prod \
--dispatch-rate-period 1s \
--msg-dispatch-rate 1000 \
--byte-dispatch-rate 1048576
Topic 类型分为以下几种:
- Persistent Topic(默认):所有消息持久化到 BookKeeper,保证数据不丢失。
- Non-Persistent Topic:消息仅存在于 Broker 内存,Broker 宕机或订阅者断开后消息即丢失,但吞吐更高、延迟更低。
- Partitioned Topic:将 Topic 拆分为多个 Partition,分散到不同 Broker,实现水平扩展与并行消费。
- System Topic:内置 Topic,用于事务协调、Schema 存储等。
Partition 的计算方式为 TopicName.getPartitionIndex(),每个 Partition 是一个独立的 Ledger,独立写入独立的 Bookie 节点组。
四、四种订阅模式与消费语义
Pulsar 的消费模型基于订阅(Subscription),而 Kafka 的消费者组模型仅对应其中一种。Pulsar 提供四种订阅模式,灵活适配不同的业务场景:
| 订阅模式 | 语义描述 | 典型场景 |
|---|---|---|
| Exclusive | 同一订阅内只有一个消费者能消费,保证消息顺序 | 需要严格全局有序(如数据库 Binlog) |
| Failover | 一个主消费者失败后,另一个接管,类似主备切换 | 高可用有序消费 |
| Shared | 多个消费者共享同一订阅,消息轮询分发 | 无序高吞吐并行消费 |
| Key_Shared | 按 Key Hash 路由到固定消费者,同一 Key 有序且并行 | 需要按 Key 保序的场景 |
// Key_Shared 订阅示例
Consumer consumer = pulsarClient.newConsumer()
.topic("persistent://useast/tenant/prod/orders")
.subscriptionName("order-processor")
.subscriptionType(SubscriptionType.Key_Shared)
.keySharedPolicy(KeySharedPolicy.autoSplitHashRange())
.subscribe();
关于消息确认(Acknowledgment),Pulsar 支持两种粒度:Cumulative Ack 将确认到某条消息之前的所有消息;Individual Ack 逐条确认,未被确认的消息将存储在 MarkDeletePosition 之后,直到被显式确认或订阅被删除。
五、Schema Registry:消息序列化的工程治理
Pulsar 内置 Schema Registry,支持强类型消息定义。Schema 附加在 Topic 级别,Broker 会在生产端校验消息格式,不匹配则拒绝写入。支持的 Schema 类型包括:
- Primitive Type:String、Int32、Boolean 等基础类型
- JSON Schema:基于 JSON Schema 标准的结构化类型
- Avro Schema:Apache Avro 格式,自带 Schema 演化能力
- Protobuf Native Schema:Protobuf 原生支持,高性能序列化
- Key/Value Schema:复合类型,适用于有 Key 的消息
// JSON Schema 使用示例
Schema jsonSchema = Schema.JSON(OrderEvent.class);
Producer producer = pulsarClient.newProducer(jsonSchema)
.topic("persistent://useast/tenant/prod/orders")
.create();
// Schema 演化:向后兼容修改
// Pulsar 自动检查兼容性
producer.newMessage()
.value(new OrderEvent("ORD-001", "CREATED"))
.send();
六、Transactions:端到端 Exactly-Once
Pulsar 2.8+ 正式支持分布式事务,实现了端到端的 Exactly-Once 语义。其核心设计包括:
- TC (Transaction Coordinator):事务协调器,管理事务生命周期
- TB (Transaction Log):事务状态持久化存储
- Exponential Backoff Retry:指数退避重试机制
- Multi-Hop Support:跨集群事务支持
// Go 客户端事务示例
txn, err := client.NewTransaction(uuid.NewUUID(), 30*time.Second)
if err != nil { return err }
// 生产端:向多个 Topic 发送事务消息
msgID, err := producer.Send(ctx, &pulsar.ProducerMessage{
Payload: []byte("txn-payload"),
Transaction: txn,
})
// 消费端:ACK 参与事务
consumer.AckIDCumulative(id)
// 提交事务
if err := txn.Commit(ctx); err != nil {
txn.Abort(ctx)
}
七、Geo-Replication:跨地域复制
Pulsar 的 Geo-Replication 是一种异步、断点续传、去重设计的多活容灾方案。与 Kafka MirrorMaker 相比,Pulsar Geo-Replication 是 Broker 层级的原生集成:
- 自动维护 Replication Cluster 映射关系
- 基于地层存储的偏移量跟踪,自动跳过已复制的消息
- 支持 star、full-mesh 等多种拓扑
- 支持 Namespace 级别选择性复制
# 创建 Geo-Replication 域
pulsar-admin clusters create beijing \
--url http://pulsar-bj.example.com:8080 \
--broker-url pulsar://pulsar-bj.example.com:6650
pulsar-admin clusters create shanghai \
--url http://pulsar-sh.example.com:8080 \
--broker-url pulsar://pulsar-sh.example.com:6650
# 在 Tenant 级别绑定两个集群
pulsar-admin tenants create global-tenant \
--admin-roles admin \
--allowed-clusters beijing,shanghai
# 创建跨集群 Namespace
pulsar-admin namespaces create global-tenant/global
pulsar-admin namespaces set-clusters global-tenant/global \
--clusters beijing,shanghai
八、Tiered Storage:分层存储与冷数据卸载
Pulsar 的 Tiered Storage 允许将超出保留窗口的消息自动卸载到底层对象存储(如 AWS S3、Azure Blob、GCS、HDFS)。核心的 Offloader 组件基于 Apache JClouds 实现:
- 当 Topic Backlog 超过阈值时,自动将旧 Ledger 上传至对象存储
- 读取冷数据时,根据 Ledger ID 自动从对象存储缓存加载
- 对上层应用完全透明,无需修改任何读写代码
- 支持 Bucket-based Offloading、Time-based Offloading 两种触发策略
# 启用 S3 Offloader
pulsar-admin namespaces set-offload-policies \
useast/tenant/archive \
--driver aws-s3 \
--bucket pulsar-offload-us-east-1 \
--region us-east-1 \
--maxBlockSizeInBytes 67108864 \
--offloadThresholdInBytes 1073741824 \
--offloadDeletionLagInMillis 14400000
九、生产级部署:Kubernetes 上的 Pulsar
Pulsar Helm Chart 是 Kubernetes 上部署 Pulsar 的官方推荐方案,覆盖以下核心组件:
- BookKeeper(有状态 StatefulSet,独立 PVC)
- Broker(Deployment,无状态)
- Proxy(Deployment,提供统一接入入口)
- ZooKeeper(StatefulSet + PVC)
- Functions Worker(独立 Deployment)
- Pulsar Manager(Web 管理界面)
# values.yaml 核心配置片段
broker:
replicaCount: 3
resources:
requests:
memory: "4Gi"
cpu: "2"
limits:
memory: "8Gi"
cpu: "4"
configData:
managedLedgerDefaultAckQuorum: "2"
managedLedgerDefaultWriteQuorum: "3"
managedLedgerDefaultEnsembleSize: "3"
maxUnackedMessagesPerConsumer: "10000"
dispatcherMaxReadBatchSize: "100"
dispatcherMaxReadSizeBytes: "5242880"
bookkeeper:
replicaCount: 4
volumes:
journal:
size: 20Gi
storageClass: gp3
ledgers:
size: 100Gi
storageClass: gp3
configData:
dbStorage_writeCacheMaxSizeMb: "512"
dbStorage_readAheadCacheMaxSizeMb: "256"
nettyMaxFrameSizeBytes: "5252880"
zookeeper:
replicaCount: 3
volumes:
data:
size: 10Gi
关键生产调优参数:
managedLedgerDefaultEnsembleSize:控制 Ledger 分配的 Bookie 节点数量managedLedgerDefaultWriteQuorum:同时写入的副本数managedLedgerDefaultAckQuorum:确认写入的最小成功副本数(通常设为 WQ 减 1 以平衡性能与可靠性)maxUnackedMessagesPerConsumer:限制单消费者未确认数量,防止内存溢出
十、可观测性:指标体系与告警规则
Pulsar 通过 /metrics 端点暴露 Prometheus 格式的全量监控指标,覆盖以下维度:
| 指标类别 | 核心指标 | 告警阈值参考 |
|---|---|---|
| Broker 层 | pulsar_rate_in / pulsar_rate_out | 流量突降 50% |
| BookKeeper | bookie_WRITE_BYTES / bookie_ADD_ENTRY_REQUEST | 写入延迟 P99 > 200ms |
| 消息积压 | pulsar_subscription_backlog_quota_limit | Quota 使用率 > 80% |
| 订阅延迟 | pulsar_subscription_backlag_bytes | 持续增长无收敛 |
| 消息分发 | pulsar_subscription_msg_rate_out | 消费速率远低于生产速率 |
| 存储 | pulsar_storage_size / pulsar_storage_backlog_size | 存储增长曲线异常 |
# Prometheus 告警规则示例
groups:
- name: pulsar-alerts
rules:
- alert: PulsarSubscriptionBacklogHigh
expr: |
(pulsar_subscription_backlog_size / pulsar_subscription_backlog_quota_limit) > 0.85
for: 5m
labels:
severity: critical
annotations:
summary: "Pulsar subscription backlog exceeds 85% of quota"
- alert: BookieWriteLatencyHigh
expr: histogram_quantile(0.99, rate(bookie_ADD_ENTRY_LATENCY_ms_bucket[5m])) > 200
for: 3m
labels:
severity: warning
十一、Functions 与 Pulsar SQL
Pulsar Functions 是轻量级流处理框架,继承 Stream-Runtime 模型,支持 Java、Go、Python 等多种语言。Functions 运行在 Broker 侧(或独立 Functions Worker),接收上游 Topic 消息,经过变换后写入下游 Topic,状态由内置的 State Storage(基于 BookKeeper)管理。
# Python Function 示例
import logging
class WordCountFunction:
def __init__(self):
self.word_counts = {}
def process(self, input, context):
words = input.split()
for word in words:
self.word_counts[word] = self.word_counts.get(word, 0) + 1
context.get_logger().info(f"Current counts: {self.word_counts}")
if self.word_counts:
max_word = max(self.word_counts, key=self.word_counts.get)
return f"Top word: {max_word} ({self.word_counts[max_word]})"
Pulsar SQL (Presto/Pulsar) 允许基于 Presto 查询 Topic 中过去的数据,通过 BookKeeper 的 ManagedLedger API 直接读取历史消息,适用于数据审计、回放和分析。
-- 查询最近 10 条订单
SELECT * FROM "persistent://useast/tenant/prod/orders"
ORDER BY __publish_time__ DESC
LIMIT 10;
-- 统计每小时生产速率
SELECT
date_trunc('hour', from_unixtime(__publish_time__/1000)) AS hour,
COUNT(*) AS msg_count
FROM "persistent://useast/tenant/prod/orders"
WHERE from_unixtime(__publish_time__/1000) > current_timestamp - INTERVAL '24' HOUR
GROUP BY 1;
十二、与 Kafka 的关键差异化对比
理解 Pulsar 相对于 Kafka 的差异化,有助于在合适的场景做出技术选型:
- 架构层面:Kafka 的 Broker 与存储耦合在一起,扩缩容需要迁移数据;Pulsar 的计算/存储分离天然支持弹性伸缩。
- 多租户:Kafka 仅支持 ACL 级别的粗粒度隔离;Pulsar 从 Namespace 级别提供资源配额、Policy 隔离、Geo-Replication 等完整的多租户方案。
- Topic 扩展:Kafka 的 Partition 一旦设定很难修改;Pulsar 的 Partition 可动态增加,且不影响已有消息。
- 订阅模型:Kafka 仅支持 Consumer Group 一种模型;Pulsar 提供四种订阅模式,更灵活。
- 运维层面:Pulsar 的 BookKeeper Compaction 自动清理过期数据,Kafka 需要依赖 Log Compaction。
当然,Pulsar 也有其局限性:社区与生态成熟度不及 Kafka;ZooKeeper 引入额外组件;小规模场景下的运维复杂度更高。因此,对于 Kafka 已经满足需求的场景,迁移 Pulsar 需要充分评估增量收益。
十三、总结与展望
Apache Pulsar 凭借其云原生架构、计算存储分离、多租户原生支持、Geo-Replication、Tiered Storage 等能力,为大规模分布式消息流场景提供了坚实的底座。随着 Pulsar 3.0 的发布,系统进一步强化了分层存储自动化、Kubernetes 原生调度、事务性能优化等方向。未来,Pulsar 在物联网数据接入、金融多活容灾、实时数仓 CDC 等领域的应用将持续深化。
掌握 Pulsar 不仅仅是学习一个中间件,更是理解现代流处理体系的机会。希望本文的架构原理与工程实践能帮助读者在生产中构建可靠、可扩展的实时数据管道。

发表评论 取消回复