楼层: 首页/ 软件技术/ 中间件全景/ 消息队列补全:Apache Pulsar
18

消息队列补全:Apache Pulsar

Message Queue · Pulsar vs Kafka vs RocketMQ

第 3 章我们讲过 Kafka 和 RabbitMQ,第 6 章提到过 RocketMQ 的事务消息。但消息队列的世界里还有一个名字绕不开:Apache Pulsar。它是"后 Kafka 时代"最被认真讨论的设计——把存储和计算彻底分开,Broker 变得无状态,数据交给专门的一层来存。这一章我们先说清"Kafka 到底哪三件事不顺手",再看 Pulsar 是怎么回答的,最后给你一张 Kafka / RocketMQ / Pulsar 的三方选型表。看完你会知道:Pulsar 不是"更好的 Kafka",它是为另一种需求形态设计的。

17.1 先看 Kafka 的三处不顺手

Kafka 的三个"设计代价"(不是缺陷,是取舍)
痛点为什么会这样实际影响
存算一体,扩容要搬数据Broker 既负责网络服务,也负责把数据落在自己的磁盘上;一个分区归属某个 Broker,数据不会自己跑存储不够时加机器,必须做分区重分配(rebalance),把大量数据在集群内网复制一遍——耗时长、影响性能,通常只能低峰做。而且一个分区的大小受单机磁盘限制
多租户隔离弱只有 topic/分区概念,没有租户层;资源配额机制相对简单一个共享集群里,某个业务建了 10 万个 topic 或疯狂写入,会挤占别人的资源;大规模共享集群的治理成本高,很多公司只能"一个团队一套集群"
延迟消息要自己造Kafka 没有"定时投递"的原生语义要实现"30 分钟后关单",得自己建多个延时 topic + 定时转发服务(时间轮),逻辑自己维护、可靠性自己保证

论Kafka 的消费模型还有一个容易被忽略的限制

Kafka 的"消费位点(offset)"是存在消费者组上的:一个分区在一个消费者组里只能被一个消费者消费。这带来两个后果:

① 同一份数据给多个消费者组,每个组都要从头建一套消费进度——想同时"实时处理"和"离线全量重算",得用两个不同的消费者组,各自维护 offset。数据本身不重复,但"消费进度"是重复管理的。

② 队列语义(一条消息被某一个消费者消费一次)和广播语义(每个消费者都收到)不能在同一套订阅里混用——必须靠不同的 group id 来区分。

Pulsar 的关键差别就在这里:它的消费进度(cursor)是挂在"订阅(subscription)"上的,而订阅由消费者自己创建。存储里只有一份消息,你可以同时挂:一个 Exclusive 订阅做有序处理、一个 Shared 订阅做并行任务、一个只在需要时创建的订阅来从头回放历史数据——三套进度互不干扰,且都不需要复制数据。这是 Pulsar 在"同一份数据要多种消费方式"场景下的最大优势。

17.2 Pulsar 的架构:存算分离是怎么做的

Producer / Consumer
用 Pulsar 协议(也有 Kafka 协议兼容层)连接 Broker
→
Broker(无状态)
负责协议解析、消息分发、订阅管理、限流与配额。不存数据,可以随意增删
→
BookKeeper(Bookie)
真正存数据的层。一条消息写到多个 Bookie,用 quorum 保证副本
元数据层
ZooKeeper 或 etcd(新版本支持):存租户/命名空间/topic 元数据、Broker 归属
→
分层存储 Tiered Storage
老数据自动卸载到 S3 / HDFS / 对象存储,Broker 只保留热数据
→
效果
存储不够加 Bookie(不用搬数据);吞吐不够加 Broker(不用搬数据)
存算分离带来的实际差别
操作Kafka(存算一体)Pulsar(存算分离)
扩容存储加 Broker + 分区重分配(搬数据)加 Bookie 节点,BookKeeper 自动把新写入分散到新节点,老数据不用搬
扩容吞吐加 Broker + 重分配分区(搬数据)加 Broker,把 topic 的归属(ownership)迁过去——只迁元数据,不迁数据
故障恢复Broker 挂了,其上的分区 leader 要重新选举,期间该分区不可用Broker 无状态,挂了把 topic 归属转给其他 Broker 即可;数据在 BookKeeper 里由副本保证
历史数据保留受 Broker 磁盘容量限制,只能靠调短保留时间开分层存储后,保留几个月甚至永久,成本取决于对象存储
存算分离不是免费的午餐

① 多了一层 BookKeeper,运维复杂度显著上升。Bookie 需要管理、监控、扩缩容,出问题时排查链路更长(Broker 正常但 Bookie 写不进去是很典型的故障)。没有专门的运维人力,别轻易自建 Pulsar。

② 写放大与延迟。一条消息要写到多个 Bookie 并等待 quorum 确认,相比 Kafka 的"写本地磁盘 + 副本异步拉"(Kafka 其实也是多副本,但路径不同),Pulsar 的端到端延迟通常略高(通常是毫秒级的差距,但对极致低延迟场景要实测)。

③ 小集群下"存算分离"的优势体现不出来。如果你只有 3 台机器、数据量也不大,Kafka 的运维简单性明显更划算。Pulsar 的优势需要规模才能兑现(多租户、大 topic 数、需要长期存储、需要频繁扩容)。

④ 生态与人才。Kafka 的周边(Connect、Streams、Flink 集成、托管服务、招聘市场)都更成熟。选型时要诚实评估"出事时能不能招到人"。

17.3 多租户:tenant / namespace / topic 三层模型

这是 Pulsar 相对 Kafka 最"企业级"的一个设计。它把"一个集群服务很多团队"当成一等需求来做,而不是让每个团队各自养一套集群。

三层模型各自管什么
层级代表什么典型配置
Tenant
租户
一个团队 / 一条业务线 / 一个客户。是隔离与计费的顶层单位谁能管理它(admin roles)、允许创建哪些命名空间、是否允许跨租户访问
Namespace
命名空间
租户下的一组 topic(相当于"某个环境的某个业务域",如 shop/order-prod)配额(发布/订阅速率、带宽、存储上限)、消息 TTL、保留策略(retention)、隔离策略(绑到特定 Broker 组)
Topic真正的消息通道,支持分区(partitioned topic)分区数、是否持久化、消息去重、压缩方式

用 pulsar-admin 把一个租户建起来,并给它加上配额和物理隔离

# 1) 建租户(并授权给某个团队的管理员) bin/pulsar-admin tenants create shop \ --admin-roles shop-admin --allowed-clusters "standalone" # 2) 建命名空间(业务域 + 环境) bin/pulsar-admin namespaces create shop/order-prod # 3) 配额:限制这个命名空间的发布/订阅速率、带宽和存储,防"邻居吵闹" bin/pulsar-admin namespaces set-publish-rate shop/order-prod \ --msg-rate 10000 --byte-rate 104857600 # 1 万条/秒、100MB/秒 bin/pulsar-admin namespaces set-backlog-quota shop/order-prod \ --limit 10G --policy producer_request_hold # 堆积超 10G 就暂停生产(而不是丢弃) # 4) 保留策略与 TTL:消息存多久(配合分层存储可以放很久) bin/pulsar-admin namespaces set-retention shop/order-prod \ --size 100G --time 30d bin/pulsar-admin namespaces set-message-ttl shop/order-prod --messageTTL 604800 # 5) 物理隔离:把重要业务绑定到专属 Broker 组(避免和其他业务抢资源) bin/pulsar-admin ns-isolation-policy set --primary 4 --secondary 2 \ shared-ns "standalone" "shop/order-prod" "broker-group-premium" # 6) 查看当前租户/命名空间清单(排障和审计都常用) bin/pulsar-admin tenants list bin/pulsar-admin namespaces list shop

论多租户真正解决的是"共享集群的信任问题"

很多公司的现实是:每个团队一个 MQ 集群,因为"怕被别人影响"。代价是集群数量爆炸、资源利用率低、运维成本高。

Pulsar 的三层模型把"共享但不互相伤害"变成了可配置的事:配额限制单个命名空间最多能用多少;保留策略限制每个命名空间占多少磁盘;隔离策略允许把核心业务钉到专属 Broker 组上;分层存储让"长期保留数据"的成本降下来(老数据进对象存储)。

所以 Pulsar 的典型用户画像是:平台团队要运营一个服务几百个内部租户的大集群。如果你的组织形态不是这样,多租户的价值就有限。

17.4 四种订阅模式:同一份数据,多种消费方式

订阅模式对比(这是 Pulsar 最有辨识度的能力)
订阅类型行为顺序性典型用途
Exclusive
独占
一个订阅同时只允许一个消费者(其他消费者连接会失败,直到当前消费者断开)严格有序需要严格按序处理的场景,如状态机、单分区账务
Failover
故障转移
多个消费者可以连,但只有一个在收消息;它挂了,其他自动接管严格有序(同一时刻只有一个消费者在消费)有序 + 高可用,生产最常用的有序方案
Shared
共享
多个消费者轮询(round-robin)消费,容量能水平扩展不保证顺序任务型消费(发短信、生成报表),关心吞吐不关心顺序
Key_Shared
按键共享
同一 key 的消息固定给同一个消费者,不同 key 可以并行同 key 有序按用户/订单维度有序,同时还能并行扩展
// 消费者侧指定订阅名与类型(订阅名必须稳定,见下面的坑) PulsarClient client = PulsarClient.builder() .serviceUrl("pulsar://pulsar-broker:6650") // TLS 用 pulsar+ssl://,默认 6651 .build(); // 队列语义 + 水平扩展:Shared 订阅 Consumer<byte[]> worker = client.newConsumer() .topic("persistent://shop/order-prod/order-events") .subscriptionName("order-worker") // ← 订阅名,生产必须固定且纳入版本管理 .subscriptionType(SubscriptionType.Shared) .ackTimeout(30, TimeUnit.SECONDS) // 30 秒没 ack 就重投(相当于"自动重试") .deadLetterPolicy(DeadLetterPolicy.builder() .maxRedeliverCount(3) .deadLetterTopic("order-events-DLQ") .build()) .subscribe(); // 按 key 有序 + 并行:Key_Shared 订阅 Consumer<OrderEvent> ordered = client.newConsumer(JSONSchema.of(OrderEvent.class)) .topic("persistent://shop/order-prod/order-events") .subscriptionName("order-by-user") .subscriptionType(SubscriptionType.Key_Shared) .keySharedPolicy(KeySharedPolicy.stickyHashRange().allowOutOfOrderDelivery(false)) .subscribe(); while (true) { Message<byte[]> msg = worker.receive(); try { handle(msg); worker.acknowledge(msg); // 处理成功才 ack(累计 ack 可减少请求数) } catch (BusinessException e) { if (isPermanent(e)) { worker.negativeAcknowledge(msg); // 明确"处理失败",尽快重投 } else { // 临时失败:不 ack,等 ackTimeout 到期自动重投 } } }
Pulsar 消费端的四个坑

① 订阅名写错 = 从零开始消费。Pulsar 的消费进度挂在订阅上。如果代码里订阅名带了随机数或改了名字,就等于创建了一个全新订阅,会从头(或从最早保留位置)消费一遍——现象是"消息突然大量重复"。订阅名要当成配置项管理,别硬编码随机值。

② Shared 订阅不保证顺序。想"并行又要有序",必须用 Key_Shared 并且同一业务实体用同一个 key(比如都用 orderId)。用错了会得到"偶尔乱序"的诡异现象,极难排查。

③ ackTimeout 机制会带来重复消费。Pulsar 客户端默认在 ackTimeout 后自动重投未 ack 的消息(不管你的消费者是否还在处理)。所以消费逻辑必须幂等——这一点和所有 MQ 一样,但 Pulsar 的重投更容易发生(比如业务处理超过 30 秒)。

④ 消息堆积(backlog)不看的话,会突然撑爆存储。消费者挂了没人消费,消息全部堆在 BookKeeper 里。必须监控 backlog size 和 subscription backlog,并设置 backlog-quota + 告警,否则就是磁盘写满的线上事故。

17.5 延迟消息与重试:Pulsar 怎么做的

"延迟投递"是业务里非常高频的需求:下单 30 分钟未支付自动关单、注册 24 小时后发提醒、失败任务 5 分钟后重试。Pulsar 原生支持任意延迟,这是它相比老版本 RocketMQ(只有 18 个固定延迟级别)的一个明显优势。

生产端:延迟消息 + 定时消息 + 消息去重

Producer<OrderEvent> producer = client.newProducer(JSONSchema.of(OrderEvent.class)) .topic("persistent://shop/order-prod/order-events") .producerName("order-api-1") // 固定名字 + 下面开去重,可实现"生产端幂等" .enableBatching(true) .batchingMaxPublishDelay(10, TimeUnit.MILLISECONDS) // 攒批 10ms,吞吐大幅提升 .create(); // 1) 延迟 30 分钟投递:下单未支付自动关单 producer.newMessage() .key(orderId) .value(new OrderEvent(orderId, "CREATED")) .deliverAfter(30, TimeUnit.MINUTES) // 相对延迟:任意时长 .send(); // 2) 指定时刻投递(绝对时间) producer.newMessage() .key(orderId) .value(new OrderEvent(orderId, "REMIND")) .deliverAt(System.currentTimeMillis() + 24 * 3600 * 1000L) .send(); // 3) 生产端去重(幂等):同一 sequenceId 重复发送,broker 只保留一条 producer.newMessage() .key(orderId) .sequenceId(nextSeqFor(orderId)) // 同一个 producerName 内单调递增 .value(event) .send();
延迟消息能力对比(选型时很实用)
产品延迟能力
Pulsar原生支持任意延迟(deliverAfter / deliverAt),精度到毫秒级(实际取决于内部检查间隔)
RocketMQ 5.x支持任意时刻的定时消息(5.0 起的定时消息基于时间轮 + 存储,精度和灵活性大幅提升)
RocketMQ 4.x只有 18 个固定延迟级别(1s/5s/10s/…/2h),不能指定任意时间——这是很多团队升级 5.x 的直接原因
Kafka没有原生延迟消息,需要自己用"多个延时 topic + 定时转发"实现
RabbitMQ可通过 TTL + 死信队列(DLX)实现延迟,或装延迟消息插件;延迟消息会先在队列里排队,量大时要注意
延迟消息的三个坑

① 延迟消息是"存在 broker 里排队"的,不是免费的。海量延迟消息会占用存储和内存(延迟索引),延迟时间越长、消息越多,堆积越大。用几百万条"延迟 7 天"的消息来做定时任务,是很糟糕的用法——那种场景该用专门的调度系统(如 XXL-Job/Quartz)配合 MQ。

② 延迟消息的精度不是绝对的。"延迟 30 分钟"通常是"不早于 30 分钟",实际投递时间取决于 broker 的检查频率和负载,可能有秒级偏差。业务上要能容忍"稍晚一点",并且靠幂等 + 状态机兜底(比如到了 30 分钟发现订单已支付,就什么都不做)。

③ 别把延迟消息当作"分布式定时任务"的唯一手段。它的定位是"消息投递的时间控制",不是"任务编排"。复杂的定时逻辑(依赖关系、失败重试策略、可视化)应该交给调度平台。

17.6 三方对比与选型:Kafka / RocketMQ / Pulsar

三大消息中间件全景对比
维度KafkaRocketMQPulsar
架构存算一体,分区归属 Broker存算一体(5.x 引入 Proxy 与分级存储,向分离演进)存算分离:无状态 Broker + BookKeeper 存储
吞吐极高(顺序写 + 零拷贝 + 批量),大数据管道标杆高(十万级/秒),延迟低高,但写路径经过 BookKeeper,延迟通常略高于 Kafka/RocketMQ
延迟消息无原生支持,需自建5.x 支持任意定时;4.x 仅 18 个固定级别原生任意延迟
顺序消息分区内有序分区内有序 + 严格顺序消息(队列级)Exclusive/Failover 严格有序,Key_Shared 按 key 有序
事务消息支持事务(但生态理解成本高)事务消息是国内最成熟的实现(半消息 + 回查)支持事务消息(无跨 topic 的通用事务)
多租户弱(靠 topic 与配额,共享集群治理成本高)中(namespace + 配额)强:tenant/namespace/topic + 配额 + 物理隔离 + 分层存储
消费模型消费进度绑在消费者组;一个分区一组内只被一个消费者消费与 Kafka 类似(消费者组 + 队列)进度绑在订阅上,同一份数据可挂多种订阅互不干扰
历史数据回放靠保留策略 + 重置 offset靠保留时间 + 重置位点订阅可指定从最早/指定时间/指定消息 id 开始,配合分层存储可长期保留
生态与人才最成熟:Connect/Streams/Flink 集成/托管服务/招聘面广国内生态好,阿里系文档与案例丰富生态和人才相对少,运维门槛高
运维复杂度中(分区重分配是主要痛点)中(NameServer 轻量,部署简单)高(多了 BookKeeper 这一层)

选四条选型判断

① 大数据管道、日志采集、流式计算(Flink/Spark 生态)→ Kafka。这是它最舒服的位置,吞吐和生态都无可替代。

② 业务消息(订单、支付、通知),需要事务消息、严格顺序、延迟消息,团队是国内 Java 技术栈 → RocketMQ。事务消息和严格顺序消息是它最扎实的能力,也是国内业务场景最需要的两项。

③ 要运营一个大共享集群,服务很多团队/租户,需要长期保留数据、频繁扩缩容 → Pulsar。存算分离 + 多租户 + 分层存储的组合,正是为这种形态设计的。但要准备好 BookKeeper 的运维投入。

④ 需求很简单(内网服务间解耦、任务分发、不需要海量堆积)→ RabbitMQ 甚至 Redis Stream 就够了。不是所有场景都需要上"重武器"——第 3 章那些轻量方案依然值得考虑。

最后一句提醒:消息队列是"引进来就送不走的中间件",它的迁移成本极高(生产者、消费者、运维、监控、告警全都要动)。选型时把"团队有没有人能运维它"放在"它有多少先进特性"前面考虑。

记
本章小结

① Pulsar 的核心是"存算分离":无状态 Broker(算)+ BookKeeper(存)+ 元数据层。扩容不用搬数据,故障转移只需迁元数据。

② 代价是 BookKeeper 这一层的运维复杂度、略高的延迟、相对小的生态与人才池。规模不够时,"存算分离"的优势兑现不了。

③ 多租户三层模型(tenant / namespace / topic)给"共享集群但不互相伤害"提供了配置手段:配额、保留策略、物理隔离、分层存储。

④ 四种订阅模式:Exclusive / Failover(有序)、Shared(并行无序)、Key_Shared(同 key 有序 + 并行)。消费进度绑在订阅上,同一份数据能给多种消费方式用,互不干扰。

⑤ 原生任意延迟消息 + 生产端去重是它的两个实用特性。注意延迟消息是堆在 broker 里的、有精度偏差、不能当调度系统用。

⑥ 选型:大数据管道用 Kafka,国内业务消息用 RocketMQ,大共享平台/多租户/长保留用 Pulsar。消费端幂等是所有方案的共同底线。