JavaDog程序狗
发布于 2026-08-14 / 2 阅读
0
0

【Pulsar】四种消费模式怎么选?

前言

🍊缘由

消费模式只有一行配置,选错却可能处处添堵

加瓦狗站在四种 Pulsar 订阅模式的分岔路口

🐣闪亮主角

大家好,我是 JavaDog 程序狗

今天聊聊 Pulsar 的消费模式,官方叫订阅类型(Subscription Type)。

创建消费者时,它看起来只是一行配置:

.subscriptionType(SubscriptionType.Shared)

但这一行决定了消息怎么分给消费者、能不能横向扩容,以及业务顺序还能不能守住。配置能随手写,业务语义可不能随手猜。

这篇不拿术语绕弯,就围绕三个问题展开:

  1. 消息是否要求顺序?
  2. 一个消费者是否扛得住?
  3. 消费者挂掉后,是否需要自动接替?

把这三个问题想明白,四种模式基本不会选错。

正文

🎯主要目标

  1. 订阅类型到底控制什么
  2. Exclusive:单消费者独占
  3. Failover:主消费、备接替
  4. Shared:多消费者分摊
  5. Key_Shared:按 key 保序并行
  6. Pulsar 4.0 对 Key_Shared 的改进
  7. 一张表完成选型

🍪目标讲解

一、订阅类型控制什么?

订阅不是 Topic 的固定属性,而是消费者读取 Topic 时采用的一套投递规则。

同一个 Topic 可以有多个订阅。不同订阅各自维护消费进度,也可以使用不同订阅类型,彼此互不影响。

Topic、订阅与消费者之间的关系

👽 人话解释

Topic 像一份报纸,订阅决定怎么读。可以一个人从头看到尾,也可以几个人分着看。报纸没变,分配方式变了。

当第一个消费者连接一个尚无消费者的订阅时,订阅类型随之确定。需要更换类型时,应先停止该订阅下的全部消费者,再按新类型重新连接。

Java 客户端的完整写法如下:

Consumer<byte[]> consumer = pulsarClient.newConsumer()
    .topic("persistent://my-tenant/my-ns/orders")
    .subscriptionName("order-processor")
    .subscriptionType(SubscriptionType.Exclusive)
    .subscribe();

四种模式的代码确实只差一行,真正需要琢磨的是这一行背后的业务要求。

Pulsar 四种订阅类型的投递结构对比

二、Exclusive:只允许一个消费者

Exclusive 是默认订阅类型。同一个 Topic、同一个订阅名下,只允许一个消费者连接;第二个消费者连接时会收到错误。

Exclusive 模式下加瓦狗独占消息处理台

.subscriptionType(SubscriptionType.Exclusive)

它适合单消费者即可处理、同时又关心消费顺序的场景,例如低吞吐的状态流转、配置变更处理等。

👽 人话解释

像公司唯一的一枚公章,同一时刻只能由一个人拿着处理。

Exclusive 的限制也很直接:不能靠增加消费者提升同一订阅的处理能力。单个消费者到达瓶颈后,需要优化消费逻辑、调整 Topic 设计,或者重新评估订阅类型。

还要注意,分区 Topic 的顺序问题不能只看订阅类型。Pulsar 官方对 Failover 明确说明为“分区内有序”;涉及分区时,本文也不把 Exclusive 简化成跨分区全局有序。业务需要全局顺序,Topic 分区方式和生产端路由同样要一起设计。

三、Failover:一个工作,其他待命

Failover 允许多个消费者连接同一订阅,但对非分区 Topic,同一时刻只有一个消费者接收消息。当前消费者断开后,其他已连接消费者会接替。

Failover 模式下主消费者工作、备用消费者待命

.subscriptionType(SubscriptionType.Failover)

👽 人话解释

像主备值班。一个人处理任务,其他人在线等着;主值班掉线,后面的人接上。

对分区 Topic,Pulsar 会把分区分配给消费者,每个分区同一时刻最多只有一个活跃消费者,因此它保证的是 分区内顺序

消费者选择规则也有区别:

  • 非分区 Topic:按消费者订阅顺序选择。
  • 分区 Topic:先看优先级,再按消费者名称的字典序排序,并尽量在最高优先级消费者之间分配分区。

Failover 适合“需要有序处理,又不能只挂一个消费实例”的场景。备用消费者会占用连接和部分资源,因此也没必要为了看起来高可用就无脑堆一排。

四、Shared:把消息分给多个消费者

Shared 允许多个消费者连接同一订阅。消息按照轮询方式分发,每条消息只交给其中一个消费者。

Shared 模式下多个加瓦狗并行处理不同消息

.subscriptionType(SubscriptionType.Shared)

👽 人话解释

像前台分快递,包裹来了就分给空闲同事。大家一起干得快,但先来的包裹未必先处理完。

Shared 的优点是容易横向扩展,适合无状态、任务之间互不依赖的处理,例如发送通知、生成缩略图、写审计日志。

代价有两个:

  • 不保证消息处理顺序。
  • 不支持累计确认(cumulative acknowledgment)。

某个消费者断开后,已投递给它但尚未确认的消息会重新分配。因此消费逻辑仍应做好幂等,别让同一条消息重投一次就多扣一次款。

五、Key_Shared:同一个 key 有序,不同 key 并行

Key_Shared 解决的是一个常见矛盾:既想增加消费者提高吞吐,又想让同一业务对象的消息保持顺序。

它的核心规则很简单:相同 key 或 orderingKey 的消息,只会分给同一个消费者处理。

加瓦狗将相同 key 的消息分到同一条有序消费通道

例如,按用户 ID 处理行为事件:不同用户可以并行处理,同一用户的事件仍落到同一个消费者。

producer.newMessage()
    .key("user-42")
    .value("用户 42 的点击事件".getBytes(StandardCharsets.UTF_8))
    .send();

Consumer<byte[]> consumer = pulsarClient.newConsumer()
    .topic("persistent://my-tenant/my-ns/clicks")
    .subscriptionName("click-aggregator")
    .subscriptionType(SubscriptionType.Key_Shared)
    .subscribe();

这个模式是否好用,很大程度取决于 key 是否选对。

  • 用户状态流转:可以用稳定的用户 ID。
  • 订单状态处理:可以用订单 ID。
  • 账户流水:可以用账户 ID。

IP、随机数、会变化的临时标识通常不是好 key。key 一变,同一业务对象就可能被路由给不同消费者,原本想保住的顺序也跟着没了。

另外,Key_Shared 保证的是同一 key 在同一时刻由一个消费者处理,不等于整个 Topic 全局有序。

六、Pulsar 4.0 改进了什么?

Pulsar 4.0 对 Key_Shared 的重点改进来自 PIP-379:Key_Shared Draining Hashes。它与 PIP-282 有关联,但不是同一项改动。

PIP-282 的交接边界与 PIP-379 Draining Hashes 的分工

先说 PIP-282。

当新消费者加入后,一部分 key 会重新映射。Pulsar 需要保证旧消费者上尚未完成的消息处理完,再让新消费者接手这些 key。PIP-282 将判断位置从原来的 readPosition 调整为 lastSentPosition,让这个交接边界更准确。

再说 PIP-379。

旧实现可能因为某个 hash 范围存在未确认消息,阻塞其他本可继续投递的消息。PIP-379 引入 draining hashes:只限制正在交接、仍有未确认消息的 hash,其他不受影响的 key 可以继续流动。

简单理解:以前更像“一处堵车,后面都等着”;现在把真正需要交接的车道单独管起来,其他车道继续走。

Pulsar 4.0 还在 Topic stats 中增加了 Key_Shared 排障指标,帮助定位被顺序约束阻塞的 key。官方发布说明提到的按 key 查询未确认消息 REST API 属于后续规划,不能当成 4.0 已经提供的能力。

这里也别把升级效果写成固定数字。实际收益取决于 key 分布、消费者变更频率、未确认消息和超时配置,应该在自己的压测与预发布环境里验证。

七、四种模式怎么选?

先看表:

模式同一订阅可连接多个消费者顺序语义适合场景
Exclusive单消费者消费;分区设计仍影响整体顺序低吞吐、单消费者、关注顺序
Failover非分区 Topic 有序;分区 Topic 分区内有序要顺序,也要消费端自动接替
Shared不保证顺序无状态任务、横向扩展
Key_Shared同一 key 有序按用户、订单、账户并行处理

再记一个不绕口的选择顺序:

  1. 不关心顺序,只想多消费者分摊:选 Shared。
  2. 关心同一 key 的顺序,还想并行:选 Key_Shared。
  3. 只需要一个消费者:选 Exclusive。
  4. 需要有序消费,并希望消费者断开后有人接替:选 Failover。

🌰 Key_Shared 避坑三连

加瓦狗检查 Key_Shared 的三个常见配置坑

1. 批量发送要按 key 组织

官方建议使用 KeyBasedBatcher,或者关闭批量。普通批量可能把不同 key 的消息放进同一个 batch,影响 Key_Shared 的投递语义。

Producer<byte[]> producer = pulsarClient.newProducer()
    .topic("persistent://my-tenant/my-ns/orders")
    .batcherBuilder(BatcherBuilder.KEY_BASED)
    .create();

2. 不要使用累计确认

Key_Shared 与 Shared 一样,不支持累计确认,应逐条确认消息。

3. key 必须稳定且分布合理

key 频繁变化会破坏业务对象的连续性;大量消息集中在少数 key 上,又会让少数消费者成为热点。上线前至少检查 key 的稳定性、基数和倾斜程度。

总结

Pulsar 四种订阅类型没有谁更高级,只有谁更符合业务:

  • Exclusive:一个消费者独占。
  • Failover:多个消费者连接,活跃消费者断开后有人接替。
  • Shared:多消费者分摊,不保证顺序。
  • Key_Shared:同 key 保序,不同 key 并行。

真正容易出问题的,不是不会写那行配置,而是没有先定义业务到底要什么顺序。把顺序范围、扩展方式和故障接替想清楚,再选模式,比出了问题再改订阅省事得多。

🍈猜你想问

如何与狗哥联系进行探讨?

加瓦狗联系方式

关注公众号【JavaDog 程序狗】,回复【入群】或【加入】,一起聊技术、聊踩坑。

你在 Pulsar 消费模式上遇到过什么问题?评论区见,狗哥陪你一起盘。


评论