Apache Pulsar 深度工程实战:分层架构与 BookKeeper 分布式日志的硬核拆解

在大消息中间件的竞技场上,Kafka 长期占据统治地位,但 Apache Pulsar 以其独特的分层架构设计正在快速崛起。Yahoo! 最初为了解决多租户、跨地域复制和存储计算耦合等问题设计了 Pulsar,如今它已成为 CNCF 顶级项目,并在字节跳动、腾讯、Apache 软件基金会内部大规模部署。本文将从源码级深入拆解 Pulsar 的核心设计思想,重点分析其分层架构、BookKeeper 集成、Geo-Replication 和 Pulsar Functions,并通过真实代码示例展示生产级部署的工程实践。

一、分层架构:存算分离的设计哲学

Kafka 的一大痛点是存储与计算耦合在一起——Broker 既负责消息路由又负责本地存储,导致扩缩容时数据迁移代价高昂。Pulsar 的解法是彻底分层:

┌────────────────────────────────────────────────┐
│         Broker (Stateless Compute Layer)        │
│   ┌─────────┐  ┌──────────┐  ┌──────────────┐  │
│   │ REST API │  │ Producer │  │  Consumer    │  │
│   └────┬────┘  └────┬─────┘  └──────┬───────┘  │
│        │            │               │           │
│   ┌────▼────────────▼───────────────▼───────┐   │
│   │           Topic / Subscription           │   │
│   │    (Managed Ledger / BK Ensembles)       │   │
│   └────────────────┬────────────────────────┘   │
┌────────────────────┼────────────────────────────┐
│    BookKeeper (Stateful Storage Layer)          │
│   ┌──────────┐  ┌──────────┐  ┌──────────┐      │
│   │ Bookie 1 │  │ Bookie 2 │  │ Bookie 3 │      │
│   │ Journal  │  │ Journal  │  │ Journal  │      │
│   │ Ledger   │  │ Ledger   │  │ Ledger   │      │
│   │ Storage  │  │ Storage  │  │ Storage  │      │
│   └──────────┘  └──────────┘  └──────────┘      │
└─────────────────────────────────────────────────┘

Broker 层无状态,可秒级扩缩容;BookKeeper 提供强一致的分布式日志存储。两者通过 Managed Ledger 抽象层进行协议交互。

核心协议交互流程

Producer 发送消息时,Broker 将请求写入 BookKeeper 的一个 Ledger(由一组 Bookie 节点组成 Ensemble),当 Ledger 写入达到大小或时间阈值时关闭并创建新 Ledger。Consumer 通过 Cursor(即订阅游标)追踪读取进度。

主题的所有 Ledger 构成一条 Managed Ledger(逻辑日志),Broker 通过 ManagedLedger 接口暴露给上层应用。订阅模型支持四种模式:

  • Exclusive:同一订阅下只允许一个 Consumer 消费,保证严格顺序
  • Failover:多个 Consumer 但仅一个活跃,主备切换
  • Shared:多个 Consumer 并发消费单条消息仅被一个 Consumer 获取
  • Key_Shared:按 Key 哈希分配给同一 Consumer,保证 Key 级别有序

二、BookKeeper:Pulsar 的存储灵魂

理解 Pulsar 必须先理解 BookKeeper。BookKeeper 核心是一个分布式预写日志(WAL)系统,其核心概念包括 Ledger(账本)、Entry(条目)、Bookie(存储节点)和 Ensemble(写副本组)。

Quorum 写入协议

BookKeeper 使用改进的 Paxos-like Quorum 协议(基于 ACK Quorum 和 Write Quorum):

class BookKeeperWriteProtocol:
    """BookKeeper Quorum 写入流程"""

    def write_entry(self, ledger_id: int, entry_id: int, data: bytes):
        ensemble = self.get_ensemble(ledger_id)  # e.g., 3 Bookies
        write_quorum = 2  # WQ=2
        ack_quorum = 2    # AQ=2

        # 1. 并行发送到所有 Bookie
        acks = []
        for bookie in ensemble:
            try:
                bookie.add_entry(ledger_id, entry_id, data)
                acks.append(bookie)
            except Exception as e:
                logger.warning(f"Bookie {bookie.id} write failed: {e}")

        # 2. 等待至少 AQ 个 ACK
        if len(acks) < ack_quorum:
            raise UnderReplicatedException(
                f"Only {len(acks)}/{ack_quorum} acks received"
            )

        return WriteResult(success=True, entry_id=entry_id)

关键参数关系:Ensemble Size (E) ≥ Write Quorum (WQ) ≥ Ack Quorum (AQ)。典型配置 E=3, WQ=3, AQ=2,意味着写 3 个节点,至少等 2 个确认即算成功。

读写分离与 Mixed Ledger

生产实践中,Pulsar 的 BookKeeper 集群可通过设置不同的 Ensemble 实现读写隔离:热数据写入高 SSD 节点组成的 Ensemble,冷数据迁移至 HDD 节点。BookKeeper 的 Data Distribution Policy 支持按 Ledger 属性选择 Bookie 节点。

此外,BookKeeper Journal 使用 Memory Mapped I/O + 顺序追加 + fsync 分组提交的方式,将随机写转换为顺序写,单节点 Journal 写入吞吐可达 100MB/s+。

三、Pulsar Topic 生命周期管理

Pulsar 中一个 Topic(或称 Namespace/Tenant 下的命名空间)由多个 Managed Ledger 构成,Ledger 是 BookKeeper 中的基本存储单元。以下是从 Broker 视角的 Topic 管理核心对象:

// broker/src/main/java/org/apache/pulsar/broker/service/persistent/PersistentTopic.java
public class PersistentTopic implements Topic {

    private final ManagedLedger ledger;  // 底层 BookKeeper 抽象
    private final ConcurrentHashMap<String, PersistentSubscription> subscriptions;

    // 消息追加入口:Producer 调用此方法
    public CompletableFuture<Void> publishMessage(
            ByteBuf data, int numMessages) {
        // 1. 构造 Entry(消息元数据 + payload)
        long ledgerId = ledger.getId();
        long entryId = ledger.getNextEntryId();

        // 2. 异步写入 BookKeeper
        CompletableFuture<Void> future = new CompletableFuture<>();
        ledger.asyncAddEntry(data, new AsyncCallbacks.AddEntryCallback() {
            @Override
            public void addComplete(Position position, ByteBuf entryData, 
                                    Object ctx) {
                // 3. 通知等待的 Consumers
                notifyConsumers(position);
                future.complete(null);
            }

            @Override
            public void addFailed(ManagedLedgerException exception, 
                                  Object ctx) {
                future.completeExceptionally(exception);
            }
        }, null);

        return future;
    }

    // 自动创建新 Ledger(轮转策略)
    private void checkLedgerRollOver() {
        if (ledger.getCurrentLedgerSize() > rolloverSizeLimit
            || System.currentTimeMillis() - ledger.getCursorsEarliestTimestamp() 
               > rolloverTimeLimit) {
            ledger.rollCurrentLedgerIfFull(timeNow);
        }
    }
}

关键在于 Ledger 轮转机制:当当前 Ledger 达到大小阈值(默认 1GB)或时间阈值(默认 4 小时),Broker 自动创建新 Ledger,并将旧 Ledger 标记为 sealed。这使得 Topic 可以被无限追加,同时 Consumer 按 Ledger 顺序读取,天然适配流式消费模式。

四、Geo-Replication:跨地域多活的核心方案

Pulsar 原生支持 Geo-Replication,无需外部组件(如 Kafka 的 MirrorMaker)。其原理是 Producer 写入本地集群后,Broker 通过内部 Replication Cluster 自动将消息推送到远端集群:

                    ┌─────────────┐
                    │  Global NS  │──── Tenant: aaa / ClusterPair
                    └──────┬──────┘
           ┌───────────────┼───────────────┐
           ▼               ▼               ▼
    ┌──────────┐    ┌──────────┐    ┌──────────┐
    │ Region-A │    │ Region-B │    │ Region-C │
    │ (Beijing)│    │(Shanghai)│    │(Silicon  │
    │          │    │          │    │  Valley) │
    └──────────┘    └──────────┘    └──────────┘

复制过程的关键实现:Source Broker 通过一个隐藏的 Replication Subscription(以 __change_events 为前缀)读取本地 Ledger 的数据,然后通过内嵌的 Replication Producer 将消息发送到目标集群。由于 Broker 是无状态的,复制逻辑完全运行在 Broker JVM 内部,无需外部进程。

开启 Geo-Replication 的方式是在 Topic 级别配置:

# 配置集群间的复制关系
bin/pulsar-admin clusters create region-b \
    --url http://region-b-broker:8080 \
    --broker-url pulsar://region-b-broker:6650

bin/pulsar-admin tenants create aaa \
    --admin-roles admin \
    --allowed-clusters region-a,region-b

# 为 Namespace 启用跨区复制
bin/pulsar-admin namespaces set-clusters aaa/ns1 \
    --clusters region-a,region-b

五、Pulsar Functions:轻量级流处理

Pulsar Functions 是一种基于 Pulsar Topic 的轻量级 Serverless 计算框架。与 Kafka Streams 或 Flink 不同,Functions 内嵌在 Broker JVM(独立进程也可)中运行,天然具备以下优势:

  • 零序列化开销:Function 处理来自 Pulsar Topic 的消息,无需反序列化即可直接操作内部格式
  • Exactly-Once Guarantee:Pulsar 的 Transaction API 保证了 Function 处理与 ACK 的原子性
  • 动态扩缩:Function 实例(Instance)可以独立扩缩,消息按 Key 路由到对应实例

以下是一个用 Rust 语言编写的 Function 示例(通过 Pulsar C++ Client SDK 绑定):

use pulsar::*;
use serde::{Serialize, Deserialize};

#[derive(Serialize, Deserialize)]
struct SensorData {
    device_id: String,
    timestamp: u64,
    temperature: f64,
    humidity: f64,
}

#[tokio::main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
    let addr = "pulsar://localhost:6650";
    let builder = Pulsar::builder(addr, TokioExecutor);
    let pulsar = builder.build().await?;

    let mut consumer: Consumer<SensorData> = pulsar
        .consumer()
        .with_topic("raw-sensor-data")
        .with_subscription_type(SubType::Shared)
        .build()
        .await?;

    let mut producer = pulsar
        .producer()
        .with_topic("aggregated-metrics")
        .build()
        .await?;

    let mut window_agg: HashMap<String, (f64, u64)> = HashMap::new();

    while let Some(msg) = consumer.next().await {
        let data = msg.deserialize()?;

        // 滑动窗口聚合(Per-Device 每分钟平均值)
        let window_key = format!("{}-{}", 
            data.device_id, 
            data.timestamp / 60_000);

        let entry = window_agg.entry(window_key)
            .or_insert((0.0, 0));
        entry.0 += data.temperature;
        entry.1 += 1;

        // 窗口末尾发送聚合结果
        if should_flush(&window_key) {
            let avg = entry.0 / entry.1 as f64;
            let metric = json!({
                "window_key": window_key,
                "avg_temperature": avg,
                "count": entry.1
            });
            producer.send(metric).await?;
        }

        consumer.ack(&msg).await?;
    }

    Ok(())
}

Pulsar Functions 状态管理

对于有状态计算(如窗口聚合、CEP),Pulsar Functions 提供了内置的状态存储 API:

// Java SDK 状态管理示例
public class WordCountFunction implements Function<String, Void> {

    @Override
    public Void process(String input, Context context) throws Exception {
        // 读写内置状态(底层由 BookKeeper Table 实现)
        for (String word : input.split(" ")) {
            long currentCount = context.getState(word.getBytes()).orElse(0L);
            currentCount += 1;
            context.putState(word.getBytes(), Longs.toByteArray(currentCount));

            // 定期将状态备份到 BookKeeper,实现故障恢复
            context.incrCounter(word, 1);
        }
        return null;
    }
}

状态存储底层是一个 BookKeeper-backed KV Store(Table Service),数据持久化在 BookKeeper 中。Function 宕机后重启时,可从 Table 恢复上次快照。

六、生产部署实战注意事项

1. BookKeeper 磁盘规划

Journal 盘必须使用 NVMe SSD(顺序写密集),Data 盘可使用大容量 HDD。Journal 盘 IOPS 决定了单 Bookie 写入 P99 延迟,建议:

  • Journal:NVMe SSD,ext4 mount with data=writeback,noatime
  • Data:HDD 或 SATA SSD,XFS mount with noatode
  • Journal 与 Data 必须物理分离,避免 I/O 争用

2. Broker JVM 调优

Broker 大量操作在堆外(Netty Direct Memory),建议配置:

# broker.conf
managedLedgerCacheSizeMB=1024          # Managed Ledger 读缓存
managedLedgerCacheEvictionWatermark=0.92
managedLedgerDefaultMarkDeleteRateLimit=5000   # Cursor ACK 速率
brokerDeduplicationEnabled=true                   # 开启精确去重
bookkeeperExplicitLacIntervalInMills=50           # LAC 同步频率

Netty Direct Memory 需要设置 -XX:MaxDirectMemorySize,建议与 Xmx 等量:-Xmx8g -XX:MaxDirectMemorySize=8g。

3. Topic Auto-Splitting 与 Partitioned Topic

Pulsar 的 Non-Partitioned Topic 通过 Ledger 轮转实现了逻辑上的无限扩展,但对于需要并行消费的场景,建议使用 Partitioned Topic。Broker 内置 auto-split 机制,可基于流量自动分裂 Partition:

# 创建 32 Partitions 的 Topic
bin/pulsar-admin topics create-partitioned-topic \
    persistent://aaa/ns1/high-throughput-topic \
    --partitions 32

# 配置自动分裂策略
bin/pulsar-admin namespaces set-auto-topic-creation aaa/ns1 \
    --enable --type partitioned --num-partitions 16

4. 延迟消息与 Pub/Sub 死信队列

Pulsar 原生支持消息延迟投递和死信队列:

// 配置延迟消息(TTL 内未 ACK 重投)
consumerBuilder
    .ackTimeout(30, TimeUnit.SECONDS)
    .ackTimeoutTickTime(5, TimeUnit.SECONDS)
    .deadLetterPolicy(DeadLetterPolicy.builder()
        .maxRedeliverCount(3)
        .deadLetterTopic("persistent://aaa/ns1/input-topic-DLQ")
        .build());

底层通过 BookKeeper 的时间索引(Timestamp Index)实现高效延迟调度,扫描效率远高于 Kafka 的 DelayQueue 轮询方案。

七、与 Kafka 的核心差异总结

维度 Apache Kafka Apache Pulsar
存储模型 Broker 本地磁盘 BookKeeper 分布式日志
扩展粒度 Partition Rebalance(数据迁移) Broker 秒级扩缩容(无状态)
多租户 无原生支持 Namespace/Tenant/ACL 原生隔离
Geo-Replication MirrorMaker(外部工具) 原生自动复制
消息消费 Consumer Group Pull Push(Broker 主动推送)
精确去重 需要外部存储(RocksDB/TI) Broker 内置去重
延迟消息 2.0+ 支持(轮询) 原生时间戳调度

Pulsar 的核心竞争力在于存算分离带来的弹性和多租户场景下的精细控制代价。如果业务场景涉及大量 Topic(10万+)、多租户隔离或跨地域复制,Pulsar 是更优选择;对于单次写入海量数据(日志收集)和极低延迟要求(<5ms P99),Kafka 仍是更稳妥的方案。

八、结语

Apache Pulsar 不是 Kafka 的简单替代品,而是面向云原生时代消息流平台的重新设计。它的分层架构虽然带来了一定的额外 Broker→BookKeeper 网络开销,但其带来的灵活性、可运维性以及原生多地域复制能力,使其在大规模消息流场景中具备独特优势。理解 BookKeeper 的 Quorum 写入协议、Pulsar 的 Ledger 轮转机制以及 Geo-Replication 的底层推模式,是高效使用 Pulsar 的关键。

在工程落地上,建议先以非核心业务(如事件广播、异步任务队列)作为切入点,熟悉 BookKeeper 调优后,再逐步迁移核心业务流。Pulsar 社区的 StreamNative 公司也提供托管服务(StreamNative Cloud),可降低自建 BookKeeper 集群的运维成本。

点赞(0) 打赏

评论列表 共有 0 条评论

暂无评论
立即
投稿

微信公众账号

微信扫一扫加关注

发表
评论
返回
顶部