☰
Apache Pulsar核心架构与实战:云原生消息队列选型指南
2026/10/9 5:57:25 网站建设 项目流程

做后端这些年,消息队列几乎是我每天都在打交道的基础设施。早些年选型基本绕不开 Kafka,直到 Pulsar 这个名字越来越频繁地出现在各种技术大会和招聘 JD 上,我才认真去把"云原生消息队列 Pulsar"这整条技术路线啃了一遍。如果你和我一样,是已经从 Kafka 或 RabbitMQ 入门、但想搞清楚 Pulsar 到底凭什么敢叫"云原生"的人,或者你是准备在下一个项目里引入消息队列、正在纠结选型的新手,这篇文章就是写给你的。我会从 Pulsar 最核心的架构思路讲起,把它的核心概念、本地实操、重复消费问题,以及和 Kafka 的选型边界一次性讲明白,尽量不给概念注水。

1. 云原生消息队列的底气:先搞懂它出生的环境

想理解 Pulsar,你得先理解"云原生"这三个字放在消息队列身上意味着什么。它不是营销词汇,而是被架构逼出来的刚需。

1.1 传统消息队列在云环境里的别扭之处

传统的消息队列,比如以 Kafka 为代表的一批系统,在设计上有一个共同点:计算和存储绑定在同一台机器上。Kafka 的每个分区都有固定的 leader 和 follower,消息文件落在 broker 的本地磁盘上。这种设计在物理机能稳定待机的年代没什么问题,但到了容器化、Kubernetes 主导的时代,麻烦就来了:

  • 一个 Pod 突然挂了,Kafka 需要重新选主、把分区的数据从别的副本同步回来,这个过程耗时且不可控。
  • 计算量大了想加 broker?可以,但分区数据要重新迁移、再均衡,整个集群会有一段比较难受的窗口期。
  • 存储量大了想加磁盘?对不起,能做的要么是扩单机磁盘,要么重新做分区迁移,运维成本很高。

这些问题本质上是"数据跟着机器走"带来的。在云环境里,大家更希望"数据是数据、计算是计算":计算节点可以随时创建销毁,存储节点负责把数据踏实存住。

1.2 从 Yahoo 内部走出来的 Apache 顶级项目

Pulsar 最早是 Yahoo 内部为了支撑大规模消息场景开发的系统,2016 年开源,后来进入 Apache 基金会并成为顶级项目。它从第一天起就在尝试回答一个问题:能不能造一个消息队列,让 broker 完全不碰持久化存储?

这个问题的答案就是 Pulsar 的分层架构。它把消息队列拆成了两层:

  • Broker(计算层):负责任务分配、协议处理、权限校验、消息缓存,但不保存持久化数据。Broker 是没有状态的,这意味着你可以随时加一台、减一台,不需要搬运任何历史数据。
  • BookKeeper(存储层):负责消息的持久化,由一组叫 Bookie 的节点组成。消息只要写进去,就有多个副本保护,存多少数据、扩多少容量,全部由这一层决定。

这就是 Pulsar 敢喊"云原生"的底气:它在架构上就是为了"任意扩展、动态调度、故障隔离"准备的,而不是把旧架构搬到容器里强行云原生。

我个人的理解是,Kafka 更像一个"前店后厂"的模式,每个 broker 既是服务柜台又是仓库;Pulsar 则像是把仓库单独拿出来,柜台只负责接待顾客。柜台的机器坏了就换一台,顾客压根感觉不到仓库发生了什么。

2. Pulsar 最核心的底牌:存储计算分离与 Segment 机制

上一节说了分层架构,这一节必须深入一点,不然你不知道 Pulsar 那些"神话"是从哪来的。

2.1 Broker 与 BookKeeper 各自管什么

先看消息的写入路径。一个生产者把消息发到 Pulsar:

  • 生产者通过 TCP 连接到某个 Broker;
  • Broker 对消息做基本校验和路由,然后把它交给 BookKeeper;
  • BookKeeper 在多个 Bookie 节点上写入多个副本(默认通常 3 份),确认完成后 Broker 才给生产者返回成功。

这个链路里,Broker 不碰磁盘,只做转发和协调。消息的持久化、副本管理全部由 BookKeeper 负责。

有意思的是,Pulsar 的存储层没有采用"一个分区一个目录"的简单做法,而是把每个分区的数据继续切成了一个个Segment(分段)。你可以把 Segment 理解为连续写入的一组消息,长度固定或时间到了就封口,然后开启新的 Segment。这些 Segment 不是死守在某个 Bookie 上的,它们会按照分配策略散落存储到不同的 Bookie 节点上。

正因为数据被切成段打散存储,Pulsar 可以得到几个很宝贵的能力:

  • 存储与计算独立扩容:计算压力大就加 Broker,存储压力大就加 Bookie,两者互不干扰。
  • 故障恢复更快:某个 Broker 挂了,Topic 可以迅速重新调度到其他 Broker 上,因为 Broker 不需要"继承"任何本地数据。
  • 老数据可以自动降冷:那些写完很久、几乎没有读流量的 Segment,可以无缝转存到对象存储(比如 S3、OSS)上,Kafka 想干这件事很难,Pulsar 天生就支持。

2.2 和 Kafka 的架构差异对照

这一段适合放到一张表里看。我自己在做选型评审时,经常用下面这个表格向团队解释两套系统的不同:

对比项PulsarKafka
存储方式计算与存储分离,数据存 BookKeeper数据存 broker 本地磁盘
Broker 状态无状态,随时扩缩有状态,与分区数据绑定
分区迁移Topic 可在任意 broker 间调度,无需搬运数据分区迁移要重新拷贝数据
扩容方向计算加 Broker,存储加 Bookie通常要加节点并做数据重均衡
多租户内置租户、命名空间、权限、配额需要额外设计
跨地域复制内置,可配置异步复制需要 MirrorMaker 等工具
数据降冷Segment 可自动 offload 到对象存储需要自研或第三方方案

看到这个表,应该能明白为什么 Pulsar 被认为是"云原生"更彻底的那一个。但也要注意,架构红利不是免费的,代价就是系统组件更多、运维门槛更高,这个我在后面第 6 节会细讲。

3. 五个核心概念一次讲透:Topic、Partition、Subscription、Cursor 与分层存储

很多初学者看 Pulsar 文档时会发现它和 Kafka 的概念有重叠但又不一样,这里我用对比的方式帮你把这些概念钉死。

3.1 Topic 与 Partition:逻辑通道与物理分片

在 Kafka 里,Topic 是我们最熟悉的概念,一个 Topic 下有多个 Partition。Pulsar 同样有 Topic 和 Partition,而且逻辑上几乎没有差别:消息发到 Topic,Topic 内部按分区并行处理,同一个分区内的消息保持顺序。

但 Pulsar 的 Partition 在存储侧的表现不一样。Kafka 的 Partition 是"一整块日志文件";Pulsar 的每个 Partition 则是由一串 Segment 组成的逻辑序列。每个 Segment 的大小是有上限的,写满自动滚动到下一个。这带来的隐藏好处是,某个 Partition 的历史数据可以被拆成多个 Segment,分布在不同的 Bookie 上,而不是死磕一台机器的磁盘。

另外,Pulsar 里的 Topic 名称是完整带路径的,例如:

persistent://public/default/my-topic

这个路径里的public是租户,default是命名空间,后面才是真正的主题名。这个设计直接支撑了多租户:不同团队可以用同一套集群,但配额、权限、存储策略完全隔离。

3.2 四种订阅模式:消息该发给谁,怎么发

Pulsar 中,Consumer 要消费某个 Topic,必须属于某个Subscription(订阅)。Pulsar 内置了四种订阅模式,这也是它比 Kafka 更"队列化"的体现:

订阅模式消费者数量顺序性适用场景
Exclusive(独占)只能有 1 个严格有序强顺序消费,如订单流水
Failover(灾备)可以有多个,同一时刻只有 1 个在工作严格有序需要高可用但不希望并行消费
Shared(共享)多个消费者同时消费不保证全局有序高吞吐、允许乱序的批量消息
Key_Shared(按键共享)多个消费者,相同 key 路由到同一消费者按 key 有序既想并发,又要同一用户/订单有序

选订阅模式是你用 Pulsar 时第一个要做的关键决策,它直接决定了系统的行为特征。比如你用了 Shared,就不能指望全局有序;你要顺序又想多消费者分摊,那就得用 Key_Shared,并且生产端要根据业务 key(比如用户 ID)发消息。

3.3 Cursor 与 ACK:消费进度到底存在哪

每个 Subscription 都有一个独立的Cursor(游标),它相当于"我读到了哪、哪些消息还没确认"。Kafka 里也有 consumer offset,概念类似。

Pulsar 的 ACK 比 Kafka 更精细一点。它支持两种确认方式:

  • 单条 ACK:确认某一条消息。
  • 批量确认(Cumulative ACK):一次性确认到某条消息为止之前的所有消息,效率更高,但只在 Exclusive 和 Failover 模式下可用。

消费完的消息不会立刻从磁盘删除,Pulsar 会根据保留策略(retention)决定历史消息存多久。而未消费的消息累计在订阅下面,叫做Backlog,可以简单理解成"欠账"。Pulsar 会限制 backlog 的额度,如果积压超过阈值,可以让生产者暂停发送或丢弃旧消息,这一步配置好了,就能避免"消费者挂了,消息无限堆积把磁盘写爆"的事故。

另外一个很重要的点是负向 ACK 与死信。你在消费时报错,不一定要立刻自动重试。Pulsar 允许你通过negativeAcknowledge(msg)表示"这条消息我没处理好",它会稍后重新投递;如果重试次数超限,消息会被丢进死信 Topic(DLQ)。这套机制让你的消费逻辑可以放心说"不行",而不怕消息阻塞卡死整个 Partition。

4. 半小时本地跑通:用 Docker 快速体验完整链路

概念说多了容易飘,还是上手来一遍最实在。这里我用 Docker 起一个 standalone 模式的 Pulsar 实例,然后用命令行和 Java 客户端各收发一次消息。

4.1 启动一个本地单机实例

如果你机器上装了 Docker,下面的命令就够了:

docker run -d --name pulsar \ -p 6650:6650 \ -p 8080:8080 \ apachepulsar/pulsar:3.1.0 \ bin/pulsar standalone

端口说明:6650是客户端生产/消费消息的 TCP 端口,8080是 HTTP 管理接口端口。等日志里出现 "Standalone started" 或者容器状态稳定后,就算启动成功了。

注意,standalone 模式只适合本地学习,绝对不要拿去生产。它把 broker、bookie 等一整套东西塞进同一个进程,并没用真正发挥 Pulsar 的架构优势。

4.2 用命令行快速验证收发消息

先通过pulsar-admin看一眼当前集群里有什么:

docker exec -it pulsar bin/pulsar-admin tenants list docker exec -it pulsar bin/pulsar-admin namespaces list public

正常情况下会看到public租户和default命名空间。然后我们往my-topic发三条消息:

docker exec -it pulsar bin/pulsar-client produce \ persistent://public/default/my-topic \ --messages "hello-pulsar" -n 3

再开一个终端消费这几条消息:

docker exec -it pulsar bin/pulsar-client consume \ persistent://public/default/my-topic \ --subscription-name my-sub \ --num-messages 3

看到消息被打印出来,这套链路就算跑通了。除了收发消息,pulsar-admin topics list可以查看当前所有 Topic,pulsar-admin topics stats可以看到生产消费速率、backlog 等状态,这些都是排查问题很常用的命令。

4.3 用 Java 客户端在代码里收发消息

命令行只是验证环境,真正写代码才是日常。我用 Java 客户端演示一下最核心的代码。

先在pom.xml里引入依赖:

<dependency> <groupId>org.apache.pulsar</groupId> <artifactId>pulsar-client</artifactId> <version>3.1.0</version> </dependency>

生产端代码:

PulsarClient client = PulsarClient.builder() .serviceUrl("pulsar://localhost:6650") .build(); Producer<String> producer = client.newProducer(Schema.STRING) .topic("persistent://public/default/my-topic") .create(); producer.send("hello pulsar"); producer.close(); client.close();

消费端代码:

PulsarClient client = PulsarClient.builder() .serviceUrl("pulsar://localhost:6650") .build(); Consumer<String> consumer = client.newConsumer(Schema.STRING) .topic("persistent://public/default/my-topic") .subscriptionName("my-sub") .ackTimeout(30, TimeUnit.SECONDS) .subscribe(); while (true) { Message<String> msg = consumer.receive(3, TimeUnit.SECONDS); if (msg == null) { continue; } System.out.println("收到: " + msg.getValue()); consumer.acknowledge(msg); }

这里有个新手特别容易踩的坑:第一次订阅已经存在的 Topic,会默认只消费最新消息。因为 Pulsar 默认的subscriptionInitialPosition是Latest,历史消息不会给你重放。如果你希望从最早的消息开始消费,需要这样设置:

.subscriptionInitialPosition(SubscriptionInitialPosition.Earliest)

我在第一次写 demo 时就是因为没注意这个参数,明明生产端发了好几条消息,消费端却什么都收不到,排查了半天才发现是这个默认值在作怪。

5. 绕不开的重复消费问题:为什么 Pulsar 也会重复投递

搜索引擎里把"消息队列重复消费问题"归到热门词不是没有原因的。很多刚接触消息队列的人都以为换一个 MQ 就能解决重复消费,实际上这是不可能的。Pulsar 的默认投递语义是at-least-once(至少一次),也就是说,消息几乎必然存在重复投递的可能。

5.1 重复消费是怎么产生的

Pulsar 里的重复消费,绝大多数来自下面几个场景:

  • ACK 超时重新投递:你在receive()之后处理消息耗时太长,超过了ackTimeout,Broker 认为这条消息没被确认,于是重新投递。
  • 消费端崩溃:你处理完消息、但还没来得及 ACK,进程崩溃了。Broker 当然认为消息没被消费,等消费者恢复后,会重新把消息发出来。
  • Shared 模式下的负载转移:某个消费者挂了,它手上的消息会被分给其他消费者;但原来那个消费者可能已经处理了一部分,这些消息就重复了。
  • 生产端重试:网络抖动导致发送超时,生产端重试会导致同一条业务消息被写入多次。

所以你看,重复消费不是 Pulsar 的缺陷,而是分布式系统的常态。只要"确认"和"处理"不是同一个原子操作,重复就永远有机会发生。

5.2 对付重复消费的标准姿势

既然不可避免,我们唯一能做的就是让消费端具备幂等性。所谓幂等,就是"同样一条消息处理十次,结果和处理一次一样"。实操中一般用这几种方案:

  • 业务唯一 ID 去重:每条消息带上业务主键,消费端把主键存入数据库唯一索引。插入重复主键时直接报冲突跳过,天然幂等。
  • Redis 去重:消费前用SETNX判断消息 ID 是否处理过,处理完写入并设置过期时间。适合高吞吐、允许短暂重复的场景。
  • 状态机校验:如果消息处理是一个流程中的一环,先查询当前状态,只有符合前置状态才继续执行。
  • 数据库更新用绝对值:比如"把账户余额设为某值",而不是"在当前值上加 10"。前者即使重放多次结果也一致。

5.3 Pulsar 侧可以减少重复的相关配置

除了在下游做幂等,Pulsar 本身也提供了一些控制手段,但不能完全消除重复:

  • 合理设置 ackTimeout:不要设得太短,给消息处理留足时间;但也不要设得太长,否则消息卡住后不能及时转移。一般是 30 秒到几分钟,按业务耗时来定。
  • 使用重试与死信机制:消费失败先negativeAcknowledge延时重试,而不是让它无限占着 backlog。重试次数到了就丢进 DLQ,至少保证主链路不阻塞。
  • 调整acknowledgementGroupTime:让 ACK 批量发送,降低网络开销,也降低"处理成功但 ACK 没发出去"的概率窗口。

这里我想说一句实在话:不要为了追求"恰好一次"去折腾一个根本做不到的机制,真正该做的是在下游设计好幂等。Pulsar 即使提供了事务等高级功能,那也只是缩小了重复范围,而不是消灭重复。生产上见过太多团队在 MQ 层反复调参,最后发现还是业务侧加唯一索引最管用。

6. 选型建议与使用体会:什么时候该上 Pulsar,什么时候留在 Kafka

最后聊聊选型。每次我写 Pulsar 的分享,必有人问"那 Kafka 是不是要完?"。我的回答是:不会,而且很多场景继续用 Kafka 完全合理。

6.1 适合上 Pulsar 的信号

如果你遇到下面这些情况,Pulsar 值得认真评估:

  • 公司已经有成熟的 Kubernetes 平台,希望中间件也能弹性伸缩。Pulsar 的 Broker 无状态特性与容器平台配合得很好,扩缩容基本没有数据迁移负担。
  • 消息体量大、延迟敏感,且不想把存储和计算绑死死。比如团队计划把历史消息长期保留,又不想让历史数据挤占实时处理节点,Pulsar 的 Segment 分层存储能直接把旧数据卸载到对象存储。
  • 多团队共用一套集群。Pulsar 内置的租户、命名空间、配额、权限体系,比自己在 Kafka 上做一层封装省事得多。
  • 需要同时覆盖"队列模型"和"流模型"。Pulsar 的 Shared/Key_Shared 订阅模式天然支持任务分发,比 Kafka 搞 consumer group 还要做均衡策略更自然。

6.2 建议继续留在 Kafka 的场景

  • 团队已经非常熟悉 Kafka,业务跑得也好好的。架构没有非换不可的理由时,换中间件是最大的浪费。
  • 集群规模不大,不想引入 BookKeeper 这套额外的存储组件。Pulsar 的组件多、监控面和运维面都更广,小团队成本不低。
  • 生态依赖重。如果你重度使用 Flink、Spark、各种 Kafka Connect 组件,Kafka 的生态还是要领先不少,连接器和资料都更丰富。
  • 纯流式管道,没有队列分发需求。Kafka 在流式数据处理上的简洁和稳定性,经过了海量场景的验证,没必要追求"新"。

6.3 我个人在实际项目中的体会

我在生产环境真正依赖 Pulsar 是在一个跨地域多租户项目里。当时最打动我的不是它转发消息有多快,而是我可以放心地让某个租户的消息涨到另一个租户的十倍而互不干扰,也可以在流量突增时毫不犹豫地加几个 Broker 容器。这种"拆得很开"带来的底气,是传统耦合架构给不了的。

但我也得坦白说,Pulsar 的运维复杂度比单层架构高不少。BookKeeper 的 Journal、Storage 要分开挂盘,Bookie 节点的磁盘 IO 和 JVM 参数都要认真调,踩过的坑和 Kafka 的运维坑种类完全不同。所以如果不是被多租户、弹性扩缩容这些需求逼到墙角,单纯想换一个更"新"的组件,我没必要折腾。

最后给正准备入门的朋友一个建议:先在 Docker 里跑通本地 standalone,把 Topic、Subscription、ACK 这几个概念亲手验证一遍,然后设计一个带唯一主键的削峰场景练手。等你能解释清楚"为什么 Pulsar 会重复投递、我又是在哪一层把它兜住的",对这套系统的理解就算真正入巷了。

需要专业的网站建设服务?

联系我们获取免费的网站建设咨询和方案报价,让我们帮助您实现业务目标

立即咨询