Apache Pulsar 分层架构与多租户消息系统设计深度实战

引言:为什么选择 Pulsar 而非 Kafka

在云原生消息系统的选型中,Apache Kafka 长期占据统治地位,但 Apache Pulsar 凭借其独特的计算存储分离架构和原生多租户支持,正在成为新一代消息基础设施的首选方案。Yahoo! 最初设计 Pulsar 是为了解决 Messaging、Streaming、存储三合一的需求,如今它已被腾讯、Splunk、Verizon 等大规模生产环境采用。

Pulsar 与 Kafka 最本质的差异在于:Kafka 的 Broker 同时承担计算(消息路由、消费协调)和存储(日志段管理)职责,扩缩容必然伴随数据迁移;Pulsar 将存储层剥离到 Apache BookKeeper,Broker 成为无状态计算节点,实现了真正的弹性伸缩。

一、分层架构解析:计算与存储分离

1.1 三层服务模型

Pulsar 由三个独立组件构成:

  • Broker(无状态计算层):负责 HTTP 请求处理、消息路由、负载均衡、Schema 校验、消息去重等。因为无状态,扩缩容只需增减 Pod。
  • BookKeeper(持久化存储层):提供分布式预写日志(WAL)服务,每个 Ledger 是一个追加写入的日志段。
  • ZooKeeper(元数据存储):管理集群归属信息、Topic 所有权、配置项等元数据。

三层独立部署时,Broker Bookie 可以按需独立扩缩容。在 Kubernetes 中,Broker 通过 StatefulSet 部署,Bookie 使用本地 PV 或本地 NVMe 盘获得最佳 I/O 性能。

1.2 Topic 命名空间模型

Pulsar 采用 persistent://tenant/namespace/topic 的全路径命名,这是多租户的物理隔离基础:

persistent://my-company/order-events/checkout-tenant-1

每个 Namespace 可以独立配置:保留策略、消息 TTL、Backlog Quota、Dispatch Rate、消息 Schema 强制策略。Tenant 则实现最高层级的资源隔离,不同 Tenant 的 Namespace 完全不可见。

二、BookKeeper 存储引擎深度解析

2.1 Ledger、Entry、Write-Ahead Log

BookKeeper 存储模型的核心抽象:

  • Ledger:只追加写入的原子日志单元,一旦 Close 即不可变。
  • Entry:Ledger 中的每条记录,包含 Entry ID、Data、LAC(Last Add Confirmed)。
  • Journal:BookKeeper 自身的 WAL,保证写入 Ledger 的事务性。

每条消息写入时,Broker 将消息封装为 Entry 追加到当前活跃的 Ledger。当 Ledger 达到大小或时间阈值后,Broker 关闭旧 Ledger、创建新 Ledger,实现日志段轮转。

2.2 写入流程与 Ensemble 机制

Client 写入 Ledger 时,关键流程:

  1. Client 向 Bookie 组(Ensemble)写入 Entry,Ensemble 大小 = Write Quorum(例如 3)
  2. 等待 Ack Quorum(例如 2)个 Bookie 确认写入
  3. 返回 LAC 更新,消费者可安全读取

Ensemble 选择策略决定了数据安全性和写入延迟:

// 使用 RackAwareEnsemblePlacementPolicy 实现机架感知
BookKeeper bk = new BookKeeper(zkConnectionString,
    new BookKeeperConfig()
        .setEnsemblePlacementPolicy(RackAwareEnsemblePlacementPolicy.class)
        .setEnsembleSize(3)
        .setWriteQuorumSize(3)
        .setAckQuorumSize(2)
);

Write Quorum = 3, Ack Quorum = 2 容忍 1 个 Bookie 故障不丢数据;Ensemble Size = 5, WQ = 3, AQ = 2 则容忍 2 个 Bookie 故障。

2.3 Journals 与 Entry Log 分离

BookKeeper 的高性能源于 Journal File 和 Entry Log File 的分离设计:

  • Journal:使用 fsync 顺序追加写入,保证事务持久性。写入 Journal 即返回客户端成功。
  • **Entry L随机读取性能。

将 Journal 部署在独立的高性能 SSD(甚至 Optane)上,Entry Log 使用大容量 HDD,可以获得最佳性价比。

三、多租户机制与资源隔离

3.1 租户管理实战

Pulsar 的 Tenant、Namespace 两级资源模型,天然适配企业级多租户需求。通过 Admin API 可动态管理:

# 创建 Tenant 并指定可访问的集群
pulsar-admin tenants create my-tenant \
  --admin-roles admin \
  --allowed-clusters us-east,us-west

# 创建 Namespace 并配置保留策略
pulsar-admin namespaces create my-tenant/my-namespace \
  --bundles 64

# 配置保留策略:保留 100GB 或 7 天
pulsar-admin namespaces set-retention my-tenant/my-namespace \
  --size 100G \
  --time 7d

# 配置 Backlog Quota:超出 10GB 则丢弃最旧消息
pulsar-admin namespaces set-backlog-quota my-tenant/my-namespace \
  --limit 10G \
  --limitTime 3600 \
  --policy producer_request_hold

3.2 Bundle 分片与负载均衡

Namespace 被划分为若干个 Bundle(默认 64 个),每个 Bundle 是负载均衡的最小单位。当 Broker 负载过高时,Namespace Bundles 会分裂并迁移到其他 Broker:

# 手动触发 Namespace Bundle 分裂
pulsar-admin namespaces split-bundle my-tenant/my-namespace \
  --bundle 0x00000000_0x80000000 \
  --unload

# 设置自动分裂阈值
pulsar-admin namespaces set-broker-bundle-data my-tenant/my-namespace \
  --bundle 0x00000000 \
  --topics 1000 \
  --msgRateIn 50000 \
  --msgRateOut 200000 \
  --bandwidthIn 100000000 \
  --bandwidthOut 300000000 \
  --memory 400000000

3.3 授权与认证集成

Pulsar 原生支持 JWT、mTLS、Kerberos 等多种认证方案,并可按 Namespace 细粒度授权:

# 授予 subscribe 权限
pulsar-admin namespaces grant-permission my-tenant/my-namespace \
  --role consumer-service \
  --actions subscribe

# 授予 produce 权限
pulsar-admin namespaces grant-permission my-tenant/my-namespace \
  --role producer-service \
  --actions produce

# 查看权限列表
pulsar-admin namespaces permissions my-tenant/my-namespace

四、消息分发模式与消费组

4.1 四种订阅类型

订阅类型 消息分发 典型场景
Exclusive 单消费者独占,保证顺序 主备切换、全局有序消费
Failover 一主多备,主挂了切换备用 高可用有序消费
Shared 多消费者随机消费,不保证顺序 高吞吐并行消费
Key_Shared 相同 Key 路由到同一消费者 按 Key 并行且有序
// Shared 订阅:多 Worker 并行消费
Consumer<byte[]> consumer = pulsarClient.newConsumer()
    .topic("persistent://tenant/ns/events")
    .subscriptionName("worker-pool")
    .subscriptionType(SubscriptionType.Shared)
    .subscribe();

// Key_Shared 订阅:订单按 orderId 并行有序
Consumer<byte[]> orderConsumer = pulsarClient.newConsumer()
    .topic("persistent://tenant/ns/orders")
    .subscriptionName("order-processor")
    .subscriptionType(SubscriptionType.Key_Shared)
    .keySharedPolicy(KeySharedPolicy.autoSplitHashRange())
    .subscribe();

4.2 消息确认与重递

Pulsar 支持消息级别(Individual)和游标级别(Cumulative)两种确认模式。消费失败重递有两种策略:

// 定时重递(指数退避)
consumer.negativeAcknowledge(messageId);
// Broker 会在 ackTimeout 后将消息重新投递

// 延迟重递(1s、5s、30s 阶梯)
consumer.reconsumeLater(message, 1, TimeUnit.SECONDS);
consumer.reconsumeLater(message, 5, TimeUnit.SECONDS);
consumer.reconsumeLater(message, 30, TimeUnit.SSETS);
// 超过最大重递次数后进入 Dead Letter Topic
consumer.reconsumeLater(message, customMetadata, 3, TimeUnit.SECONDS);

五、Topic Compaction 与分层存储

5.1 Compaction:只保留 Key 最新值

对于"配置更新"、"物化视图"、"状态快照"这类场景,只需保留每个 Key 的最新值。Pulsar Compaction 在后台合并 Ledger,删除旧版本记录:

# 开启 Topic Compaction
pulsar-admin topics compact persistent://tenant/ns/config-topic

# 配置 Compaction Threshold:超过 1GB 自动触发
pulsar-admin topics set-compaction-threshold \
  persistent://tenant/ns/config-topic \
  --threshold 1073741824

Compaction 的实现机制:Broker 创建新 Compaction Ledger,读取所有历史 Entry,仅写入每个 Key 的最新值,然后切换读指针到新 Ledger、删除旧 Ledger。

5.2 Offload 分层存储

对于需要长期保留但又允许稍高延迟的冷数据,Pulsar 支持将旧 Ledger 卸载到 S3/OSS/HDFS:

# 配置 Offloader:卸载到 AWS S3
pulsar-admin namespaces set-offload-policies my-tenant/my-namespace \
  --driver aws-s3 \
  --bucket pulsar-offload-us-east \
  --region us-east-1 \
  --maxBlockSizeInBytes 67108864 \
  --offloadThresholdInBytes 10737418240 \
  --offloadDeletionLagInMillis 14400000

# 手动触发 Offload
pulsar-admin topics offload \
  persistent://tenant/ns/history-data \
  --sizeThreshold 10G

该机制的关键优势:卸载对消费者透明,读取时自动从对象存储拉回 BookKeeper 缓存,无需修改应用代码。

六、跨地域复制与灾备

6.1 自动 Geo-Replication

Pulsar 在 Broker 层实现跨集群复制,无需外部同步工具。配置后自动同步指定 Namespace 的所有消息:

# 在目标集群配置 Remote Cluster
pulsar-admin clusters create us-west \
  --url http://pulsar-us-west.example.com:8080 \
  --broker-url pulsar://pulsar-us-west.example.com:6650

# 在 Namespace 上启用复制
pulsar-admin namespaces set-clusters my-tenant/my-namespace \
  --clusters us-east,us-west

6.2 灾备切换策略

常见的跨地域部署模式:

模式 RPO RTO 说明
Active-Standby 秒级 分钟级 平时 Standby 不服务,主故障时切流
Active-Active ~0 秒级 就近生产消费,跨域配置/消息双向同步
Multi-Site 秒级 分钟级 多个站点互为备份
# Broker 配置:限制跨域复制带宽
replicationMetricsEnabled=true
replicationProducerQueueSize=1000
replicationConnectionsPerBroker=2

七、性能调优实战

7.1 Broker 关键参数

# broker.conf
# 消息调度线程数,建议与 CPU 核心数一致
numWorkerThreads=16

# IO 线程数,默认 CPU * 2
numIOThreads=32

# HTTP 请求处理线程
numHttpServerThreads=8

# 消息发布接收队列大小
maxMessagePublishBufferSizeInMB=5

# Managed Ledger 缓存大小(控制内存中缓存多少 Ledger)
managedLedgerCacheSizeMB=1024

# Managed Ledger 轮转检查周期
managedLedgerCacheEvictionFrequency=100

# BookKeeper 客户端最小写入 Bookie 数
bookkeeperClientMinNumWrites=2

# BookKeeper 客户端写入超时
bookkeeperClientTimeoutInSeconds=30

# Dispatch 限流:每秒最大分发消息数
dispatchThrottlingRateInMsg=0  # 0 表示不限
dispatchThrottlingRateInByte=0

# 批量发送最大批次大小
maxMessageSize=5242880

7.2 Bookie 调优

# bookie.conf
# Journal 目录建议独立 SSD
journalDirectory=/mnt/nvme-journal/bookie-journal
ledgerDirectories=/mnt/ssd-storage/bookie-ledger

# Journal sync 间隔(ms),影响写入延迟
journalSyncData=true
journalAdaptiveGroupWrites=true
journalMaxGroupWaitMS=10

# 直接内存大小,建议物理内存 1/4
dbStorage_writeCacheMaxSizeMb=2048
dbStorage_readCacheMaxSizeMb=2048

# Garbage Collector 周期(秒)
isForceGCAllow=true
gcWaitTime=300000

7.3 生产环境最佳实践清单

  1. Bookkeeper Journal 必须用独立 SSD:NVMe + noatime + deadline/none scheduler。
  2. Broker 内存配置:缓存占物理内存 25%,JVM 堆分配不超过 8GB(利用堆外内存)。
  3. 监控指标:重点关注 pulsar_rate_in/out、pulsar_storage_size、pulsar_subscription_backlog、bookie_WRITE_BYTES/READ_BYTES。
  4. Namespace 默认限制:每个 Namespace 必须设置 Backlog Quota,防止单个 Topic 无限增长占满所有 Bookie 空间。
  5. 升级策略:Broker 无状态,可滚动升级;Bookie 升级需逐个进行,确保 Write Quorum 可用。

八、与 Kafka 架构对比的深层差异

维度 Kafka Pulsar
存储架构 Broker 本地存储 BookKeeper 独立存储层
扩容成本 需数据迁移 Broker 即扩即迁移
多租户 后期增加 原生设计
消费模型 Partition-based Broker 路由,Partition 透明
跨地域复制 MirrorMaker 外部工具 内置自动
分层存储 Tiered Storage(2.4+) 原生 Offload
Geo-Replication 依赖工具 原生支持

Kafka 的优势在于生态成熟度、运维工具链丰富;Pulsar 的优势在于灵活的多租户、弹性伸缩能力、云原生友好性。

九、生产部署案例:混合云消息平台

一个典型的生产部署场景:

                    ┌─────────────────────────────┐
                    │    Kubernetes Cluster       │
                    │  ┌─────────────────────────┐│
                    │  │     Pulsar Broker       ││
                    │  │  (Deployment, HPA)      ││
                    │  └─────────┬───────────────┘│
                    │            │ HTTP            │
                    │  ┌─────────▼───────────────┐│
                    │  │     Pulsar Proxy         ││
                    │  │  (Load Balancer Entry)   ││
                    │  └─────────────────────────┘│
                    └─────────────────────────────┘
                                   │
                    ┌──────────────▼──────────────┐
                    │   StatefulSet: Bookie        │
                    │  ┌─────┐ ┌─────┐ ┌─────┐    │
                    │  │BK-0 │ │BK-1 │ │BK-2 │    │
                    │  │NVMe │ │NVMe │ │NVMe │    │
                    │  └─────┘ └─────┘ └─────┘    │
                    └─────────────────────────────┘

在 Kubernetes 部署中,使用 Helm Chart 简化管理:

# 添加 Pulsar Helm 仓库
helm repo add apache https://pulsar.apache.org/charts
helm repo update

# 部署 Pulsar(含 BookKeeper 3 节点、Broker 2 节点)
helm install pulsar apache/pulsar \
  --values my-values.yaml \
  --set initialize=true

十、监控告警与可观测性

10.1 Prometheus 集成

Pulsar 原生暴露 Prometheus 指标端点(:8080/metrics),关键监控面板应包含:

# prometheus.yml 配置
scrape_configs:
  - job_name: 'pulsar-broker'
    static_configs:
      - targets:
        - pulsar-broker-0:8080
        - pulsar-broker-1:8080
        - pulsar-broker-2:8080
  - job_name: 'pulsar-bookie'
    static_configs:
      - targets:
        - pulsar-bookie-0:8080
        - pulsar-bookie-1:8080

10.2 关键告警规则

# Alibaba 推荐的告警规则示例
groups:
- name: pulsar-alerts
  rules:
  - alert: PulsarBrokerHighBacklog
    expr: pulsar_subscription_backlog > 1000000
    for: 5m
    labels:
      severity: warning
    annotations:
      summary: "Pulsar 消费积压超过 100 万条"

  - alert: PulsarBookieReadOnly
    expr: bookie_STATUS == 0
    for: 30s
    labels:
      severity: critical
    annotations:
      summary: "Bookie 进入只读模式,写入将失败"

  - alert: PulsarBrokerHighPublishLatency
    expr: histogram_quantile(0.99, rate(pulsar_entry_size_bucket{quantile="0.99"}[5m])) > 5000
    for: 3m
    labels:
      severity: warning
    annotations:
      summary: "消息发布 P99 延迟超过 5s"

总结

Apache Pulsar 通过将计算与存储分离,从根本上解决了 Kafka 在云原生场景下面临的弹性缩容难题。其原生多租户架构使大型企业能够安全地将消息基础设施作为平台即服务(PaaS)交付给数十个业务团队。Pulsar Functions(轻量级流处理)+ Pulsar IO(数据管道)+ Topic Compaction + Offload 分层存储的组合能力,使其从传统消息队列演进为完整的流数据平台。

选择 Pulsar 而非 Kafka 的核心判断标准包括:是否需要原生多租户隔离、是否频繁扩缩容、是否需要跨地域容灾、是否长期保留大规模历史数据。当这些需求满足两条以上时,Pulsar 的架构优势将显著超过其运维复杂度的增加。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部