一、引言:为什么需要 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%
BookKeeperbookie_WRITE_BYTES / bookie_ADD_ENTRY_REQUEST写入延迟 P99 > 200ms
消息积压pulsar_subscription_backlog_quota_limitQuota 使用率 > 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 不仅仅是学习一个中间件,更是理解现代流处理体系的机会。希望本文的架构原理与工程实践能帮助读者在生产中构建可靠、可扩展的实时数据管道。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部