![注释]: 这是一篇关于Pulsar的深度技术分享,预期阅读时长:演示与实操为主,全文约 20 分钟。
1. 初识Pulsar:云原生时代真正把存储和计算拆开的那个消息队列
做后端这几年,消息队列用过好几代。从最开始跟着业务瞎折腾ActiveMQ,到后来大规模落到Kafka,再到业务量涨起来之后被各种运维问题追着跑,说实话,我对“消息队列”这四个字是有一些肌肉记忆的。直到后来因为业务需要,把Pulsar引入生产环境做了一轮完整验证,才发现这个圈子并不只是Kafka以外的“另一个选择”,而是从架构思路上就把很多从前默认无解的问题,变成了可以优雅解决的方案。
如果你现在正面对这几个问题,我强烈建议花点时间看看Pulsar:
- 团队在做微服务拆分,需要一套可靠的消息中间件,但不想投入太多人力去运维ZooKeeper和Broker集群;
- 业务有明显的流量毛刺(比如每天中午、晚上的高峰),不想为峰值流量常年预留几台空闲机器;
- 对消息投递可靠性要求高,不允许丢失,同时对重复消息又很敏感,需要精细的消费确认机制;
- 或者你已经在用Kafka,但被分区扩容、Rebalance 毛刺、存储空间暴涨这些问题折腾过。
Pulsar 最核心的一句话我用一句话给你说清楚:它把“计算”和“存储”彻底分开,Broker 无状态化,数据全部落到 BookKeeper 里。这听起来好像只是架构上“挪了一下屁股”,但带来的连锁反应是:扩容变得极其简单、存储可以独立扩展、客户端接入体验大幅提升。
这篇文章我会从它的核心架构讲起,然后直接上手跑一个真实的消息发送与消费示例,最后重点聊一个大家在做消息队列时几乎都会踩的坑——重复消费问题。全文偏实践,尽量少讲虚的。
需要模型API调用? 免费领10W Token,多模型网关一键接入 Claude、DeepSeek 等主流模型。
2. 为什么说 Pulsar 的架构设计是“云原生”的
我第一次看 Pulsar 的架构图时,说实话并没有立刻体会到它的精妙之处。直到我对比着Kafka的架构部署了一遍,才真正理解“云原生”这三个字不是营销话术。
2.1 分层架构:Broker 与 BookKeeper 的职责分离
传统消息中间件(包括 Kafka)的 Broker 是既管路由、又管存储的。Kafka 的每个分区有主副本和从副本,所有副本数据都落在 Broker 本地磁盘上。这样设计的问题在于:扩容一个分区或迁移副本,本质上是搬数据;Broker 的负载和磁盘容量被绑死在了一起。
Pulsar 把这个耦合解开了。它把消息的存储层独立成 BookKeeper 集群,Broker 只负责接收请求、管理游标、做订阅分发。消息一旦被确认写入 BookKeeper,Broker 本地几乎不落任何数据。
这意味着什么?意味着你可以随时把一个 Broker 从集群中摘掉,不需要迁移任何数据;也意味着你可以在流量高峰期临时加 Broker 节点分担压力,低峰期再缩回来。这天生就是为了云环境设计的——云上最值钱的就是弹性。
2.2 存算分离带来的直接收益
我列一个对比表,你看完就明白为什么大家越来越关注 Pulsar:
| 能力 | 传统 Kafka 架构 | Pulsar 分层架构 |
|---|---|---|
| 扩容分区 | 需要做数据重分布,耗时长、有风险 | 加 Broker 即可,计算层无状态 |
| 存储扩容 | 需要为 Broker 挂新盘或迁移数据 | 横向扩 BookKeeper 节点即可 |
| 流量高峰期 | 必须提前预留资源 | 动态扩 Broker,低峰期缩容 |
| 地域复制 | 需要额外搭 MirrorMaker 等同步工具 | 内置跨地域复制,配置即用 |
| 客户端 | 只有 Producer/Consumer | Producer/Consumer/Reader 三种角色 |
| 消息保留 | 基于 offset,靠 log retention 清理 | 基于游标,支持按时间或大小保留,可无缝回溯 |
尤其是“消息保留”这一点,我不想让你觉得这只是个存储细节,它实际影响的是整个消费模型的设计。Kafka 的消息像流水账本,消费者靠 offset 指针定位;Pulsar 的消息是持久的、有独立 ID 的条目(Entry),每个消费者(实际上叫订阅 Subscription)有自己独立的光标(Cursor)。这就引出了 Pulsar 一个非常有意思的能力——同一份数据,可以被不同订阅用完全不同的速率和逻辑来回消费,互不干扰。
2.3 订阅模型:不止是点对点和发布订阅
Pulsar 有四种订阅模式:
- 独占订阅(Exclusive):一个订阅同时只允许一个消费者,适合严格有序的场景;
- 共享订阅(Shared):消息被多个消费者轮询分发,适合吞吐量大、不要求全局顺序的场景;
- 故障转移订阅(Failover):多个消费者,但只有一个活跃接收消息,宕机后自动切换;
- Key_Shared 订阅:消息按 key 哈希到不同消费者,同一 key 的消息始终由同一个消费者处理。
我最喜欢的是 Key_Shared。以前用 Kafka 想保证同一个订单 ID 的消息被有序处理,得用分区加 key 分区器,一不小心分区数变了就全乱了。Pulsar 的 Key_Shared 直接在订阅层解决,不需要关心分区关系,对业务开发来说要友好得多。
3. 快速上手:本地跑起一个 Pulsar 集群并完成收发消息
理论吃不饱,直接上实操。我建议你直接用 Docker 跑一个单机版 Pulsar,先把流程走通,再考虑部署集群或上K8s。你不需要一开始就搞懂所有配置项,重点是建立直观感受。
3.1 环境准备与容器启动
这里我特意只映射了 6650(客户端端口)和 8080(HTTP 管理端口),另外把 BookKeeper 的数据目录挂到了宿主机,防止容器重建丢数据。
bash复制docker run -d \
--name pulsar \
-p 6650:6650 \
-p 8080:8080 \
-v /data/pulsar/data:/pulsar/data \
apachepulsar/pulsar:3.3.1 \
bin/pulsar standalone
启动需要等一会儿,因为要初始化 BookKeeper 和构建元数据。看到控制台输出 messaging service is ready 类似的日志就说明准备好了。
验证一下端口和集群状态:
bash复制curl http://localhost:8080/admin/v2/clusters/standalone
返回 {"serviceUrl":"http://localhost:8080/","brokerServiceUrl":"pulsar://localhost:6650/"},就说明管理接口和消息端口都活着。
提示:如果你是苹果芯片的 Mac,
apachepulsar/pulsar镜像也有 arm64 版本,可以直接拉取。如果是老版本遇到容器起不来,多半是内存不够,Pulsar 单机版默认 JVM 堆内存比较大,可以通过环境变量调小,比如PULSAR_MEM="-Xms512m -Xmx512m"。
3.2 用命令行工具感受消息收发
启动完成后,Pulsar 自带命令行工具 pulsar-client。为了让你更直观地看到消息轨迹,我建议开两个终端。
第一个终端订阅消息:
bash复制docker exec -it pulsar bin/pulsar-client consume \
--subscription-name my-sub \
--num-messages 0 \
test-topic
第二个终端发送 5 条消息:
bash复制docker exec -it pulsar bin/pulsar-client produce \
--messages "hello-pulsar-1,hello-pulsar-2,hello-pulsar-3" \
test-topic
你会看到消费者终端打印出每条消息的 topic、分区、消息 ID(比如 3:0:-1 这种格式)和内容。这个 ID 后面我们排查重复消费时会反复看到它,先留个印象。
这里注意,我用了 --num-messages 0,意思是持续监听不退出。实际生产脚本里不会这么干,但用于体验很直观。
3.3 Java 客户端接入:一个完整的发送与消费示例
命令行只是开胃菜。真实业务场景肯定要用 SDK。Pulsar 官方对 Java 的支持最完善,Spring Boot 集成也有官方 starter。我直接给你一个最小可用的例子,不依赖 Spring,方便你理解底层机制。
先加依赖(Maven):
xml复制<dependency>
<groupId>org.apache.pulsar</groupId>
<artifactId>pulsar-client</artifactId>
<version>3.3.1</version>
</dependency>
生产者示例:
java复制import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.Producer;
import org.apache.pulsar.client.api.Schema;
public class PulsarProducerDemo {
public static void main(String[] args) throws Exception {
// 1. 创建客户端:指向 Pulsar Broker 的消息端口
PulsarClient client = PulsarClient.builder()
.serviceUrl("pulsar://localhost:6650")
.build();
// 2. 创建生产者:指定 topic 和消息类型
Producer<String> producer = client.newProducer(Schema.STRING)
.topic("persistent://public/default/order-topic")
.create();
// 3. 发送消息
for (int i = 0; i < 100; i++) {
String content = "order-" + i;
// send 是异步接口,这里调用 get() 只是为了演示同步等待结果
var messageId = producer.newMessage()
.key("order-key-" + (i % 10))
.value(content)
.send()
.get();
System.out.println("发送成功: " + content + ", messageId=" + messageId);
}
// 4. 释放资源
producer.close();
client.close();
}
}
这段代码里有一个很重要的细节值得展开说说,就是第三条里的 .key()。别小看这个 key,它直接决定了消息在 Key_Shared 订阅模式下会被哪个消费者处理。业务中如果希望通过同一个 key(比如订单号)把相关联的消息路由到同一个消费者做有状态处理,就必须在这里指定 key。
消费者示例:
java复制import org.apache.pulsar.client.api.PulsarClient;
import org.apache.pulsar.client.api.Consumer;
import org.apache.pulsar.client.api.SubscriptionType;
public class PulsarConsumerDemo {
public static void main(String[] args) throws Exception {
PulsarClient client = PulsarClient.builder()
.serviceUrl("pulsar://localhost:6650")
.build();
Consumer<String> consumer = client.newConsumer(Schema.STRING)
.topic("persistent://public/default/order-topic")
// 订阅名很关键:同一个订阅名下的消费者共享消息
.subscriptionName("order-service")
// 共享订阅,适合水平扩展消费者
.subscriptionType(SubscriptionType.Shared)
.subscribe();
while (true) {
// 同步阻塞等待消息,实际项目里一般用 listen 或者异步 receive
var message = consumer.receive(5, java.util.concurrent.TimeUnit.SECONDS);
if (message != null) {
try {
String value = message.getValue();
System.out.println("收到消息: " + value + ", messageId=" + message.getMessageId());
// 业务处理成功后,一定要确认消息
consumer.acknowledge(message);
} catch (Exception e) {
// 处理失败,不确认,并把消息标记为待重投
consumer.negativeAcknowledge(message);
}
}
}
}
}
这段消费者代码里埋了两个很容易被新手忽略的“生死线”:acknowledge 和 negativeAcknowledge。后面我会专门展开讲,这里你只需要记住一个原则:看到一条消息,先别急着确认,等你的业务逻辑真正处理完成、落库成功了,再 ack。 如果你收到消息就立刻 ack,然后处理过程中应用宕机,这条消息就永远丢失了。
实操心得:很多第一次接触 Pulsar 的同学会把 ack 类比成“签收快递”,收到就签收。这个类比在 Pulsar 里是错的。Pulsar 的 ack 更像是“我确认把这活干完了”,不是“我收到活儿了”。这个思维转变非常重要,搞清楚了,后面很多问题都一通百通。
3.4 三种客户端角色选型
前面提了 Pulsar 的客户端有三种角色,这里顺手展开一下:
- Producer:负责发消息。
- Consumer:负责订阅消费,消费进度由 Broker 端的 Subscription Cursor 管理。
- Reader:一个“从头读”的角色,不注册订阅,完全靠手动指定消息 ID 来读取。
Reader 是 Pulsar 独有的,Kafka 里面没有对应的概念。它适合什么场景?比如你要做数据灌库、离线分析、或者实现一个“从昨天 12:00 开始重新跑一遍”的补偿任务。用 Consumer 很难优雅地做到,用 Reader 一行代码就解决了。
java复制Reader<byte[]> reader = client.newReader()
.topic("persistent://public/default/order-topic")
// 指定从这条消息开始读
.startMessageId(MessageId.earliest)
.create();
这个能力在生产环境的价值很大。遇到数据修复、逻辑变更重新计算之类的场景,你不用再临时造一套消费者,直接用 Reader 读取历史消息即可。
4. 核心机制拆解:消息确认、游标与消费进度管理
弄懂 Pulsar 的消息确认机制,是理解整个 Pulsar 消费模型的地基。我在这里多花点篇幅,因为这直接关系到重复消费。
4.1 消息确认(Acknowledgment)到底是怎样工作的
先扫清一个概念:Pulsar 的消息确认是按消息 ID(MessageId) 来做的,不是按“批次”或“偏移量”来做。每条消息进入 topic 后都会获得一个全局唯一的 ID,这个 ID 由 (ledgerId, entryId, partitionIndex) 组成。
当你调用 consumer.acknowledge(message) 时,你告诉 Broker:这个 ID 对应的消息我处理完了。Broker 收到确认后,会更新订阅的游标。
那么问题来了:如果我在 Shared 订阅模式下并发处理 100 条消息,每条都得单独 ack 吗?是的,每条消息都有一个独立 ID,都需要单独确认。但 Pulsar 为了减少确认开销,做了一个优化——累积确认(Cumulative Acknowledgment)。
累积确认的机制是:在单分区单个订阅内,你如果 ack 了一个消息 ID,Broker 会认为这个 ID 之前的所有消息也都确认了。这个优化在独占或灾备订阅下非常高效。但在共享订阅下要小心,因为消息被分发到多个消费者,你先确认了某条 ID 比较大的消息,不代表 ID 小的消息也被其他消费者处理完了。所以在 Shared 模式下,Pulsar 默认使用单条确认(Individual Acknowledgment),你可以显式设置:
java复制consumer.acknowledge(message.getMessageId());
我看过不少人把两种确认模式混在一起用,结果出现“消息丢失”的假象。这里给你一个明确的操作建议:
- Exclusive / Failover 订阅:放心用累积确认,性能好,语义安全。
- Shared / Key_Shared 订阅:永远用单条确认,不要用累积确认,否则会误确认掉其他消费者还没处理的消息。
4.2 游标(Cursor)与消息保留机制
Pulsar 里每个订阅都有一个游标,记录着这个订阅已经确认到哪里了。消费者宕机恢复后,Broker 会根据游标位置继续投递未确认的消息。
这个设计有一个很实用的后果:你消费过的消息,并不会立刻被删除。只要订阅游标没越过它,消息就还在 BookKeeper 里躺着。而且你甚至可以手动把游标往回拨,重新消费历史数据。
我在生产环境就用过这个能力。有一次业务方数据算错了,希望把订单系统过去 2 小时的消息重新消费一遍。当时我已经把消息都消费完了,游标早就到头了。本来以为要重发消息,后来查了一下文档,Pulsar 支持直接重置订阅游标到指定时间点:
bash复制bin/pulsar-admin topics reset-cursor \
--subscription order-service \
--time "2h ago" \
persistent://public/default/order-topic
命令执行完,消费者会自动开始重放过去 2 小时的消息。这个功能,Kafka 也有(通过 kafka-consumer-groups --reset-offsets),但 Pulsar 的实现更直观:时间点、消息 ID、最早/最新位置任选。
4.3 保留策略(Retention Policy)和存储水位
因为消息存储和 Broker 解耦了,Pulsar 的保留策略可以配置得非常灵活。你可以按时间、按大小来设置保留:
bash复制bin/pulsar-admin topics set-retention \
--size -1 \
--time -1 \
persistent://public/default/order-topic
-1 -1 表示永久保留。有些人看到这里会担心存储无限膨胀,实际上 BookKeeper 有自动清理机制,只有当所有订阅的游标都越过某段消息,且超过保留时间后,这段数据才会被自动删除。
这种设计带来的一个隐形成本是:消费慢的订阅会拖住存储回收。如果一个订阅长期不消费,游标不动,即使你已经用另一个订阅把消息消费完了,BookKeeper 依然要保留这些数据。所以我建议你初期就为不重要的 topic 设置合理的保留时间,避免无谓的存储占用。
5. 从 Kafka 迁移到 Pulsar 前,你需要知道的几个关键差异
如果你之前是 Kafka 的重度用户,直接从 Kafka 的思维模型去理解 Pulsar,会有几个明显的冲突点。这里我把最容易踩的差异列一下,帮你少走弯路。
5.1 Partition 的概念:从“存储单位”到“并行度单位”
Kafka 中 Partition 是物理存储单位,消息分布在多个 Partition 上,Partition 的数量决定了并发上限,也决定了存储扩容的最小粒度。Partition 一旦定下来,增减操作非常麻烦。
Pulsar 中 Topic 是逻辑概念,底层数据是被切分成分片(Fragment) 存储在 BookKeeper 的 Ledger 中的。Topic 的存储容量可以远超单个 Broker 的磁盘容量。订阅的并行度由 Subscription 的消费者数量决定,和 Partition 数量没有强绑定关系。
这就带来一个实践上的差异:在 Kafka 里你可能为了并发度预先设置 64 个 Partition;在 Pulsar 里你完全可以先用一个 Topic,后续消费者变多了,并行度自然就上去了。
5.2 消费位移(Offset) vs 消息 ID(MessageId)
Kafka 里 Consumer 的 offset 是整数,Broker 端保存。Pulsar 的消息 ID 是复合结构,可以精确定位到某一条消息,而不是某个偏移量。这也让 Pulsar 的“单条消息确认”成为可能。
5.3 从消费者拉取(Pull)到 Broker 推送(Push)
Kafka 的消费者是从 Broker 拉数据(Pull)。Pulsar 虽然底层也是长轮询,但给用户的编程接口是订阅式的,Broker 主动往客户端推送消息,客户端在本地有一个接收队列。这个队列可以用 receiverQueueSize 配置:
java复制client.newConsumer()
.receiverQueueSize(1000)
.subscribe();
调大接收队列可以提升吞吐,但也会带来更多的本地堆积,如果你的业务是逐条处理且需要严格控制顺序,队列太大反而会让单条消息延迟升高。这个值的设置要根据业务取舍。
注意:在独占或灾备订阅下默认开启
receiverQueueSize。在共享订阅下,receiverQueueSize默认按消费者数量均分。调整队列大小时要结合消息处理耗时,避免“处理不过来但 Broker 还在不断推”的情况。
6. 重复消费问题:为什么消息队列重复消费是常态,以及 Pulsar 里如何应对
聊完架构和基础用法,终于轮到全网搜索热度最高的关键词——“消息队列重复消费问题”了。很多人第一次遇到重复消费时,第一反应是“消息队列是不是出 bug 了”。我在这里明确告诉你:消息队列永远做不到“恰好一次消费”,分布式系统里这是理论极限。 你只能做到“恰好一次处理”,而“恰好一次处理”的实现方式,绝不是靠消息队列本身。
6.1 重复消费是怎么产生的
在 Pulsar 中,重复消费主要由三类原因导致:
- 消费端 ack 丢失:消费者处理完消息,向 Broker 发起 ack,但网络抖动导致 ack 没有送达。Broker 认为消息未被确认,于是超时后重新投递。
- 消费端宕机:消费者处理完消息、落库了,但还没来得及 ack 就宕机了。重启后游标停留在处理前的位置,这条消息被再次投递。
- ack 超时机制(ackTimeout):Pulsar 允许你设置消息确认超时时间。如果消费者在超时时间内没有 ack,Broker 就自动把消息重新投递给其他消费者。如果你处理逻辑耗时较长,就很容易被误判。
第三种情况是新手最容易踩的。你明明在正常处理,只是处理比较慢,Broker 却把你的消息标记为超时,重新投递了。两个消费者同时处理同一条消息,业务侧出现重复写入。
6.2 解决重复消费的第一原则:业务幂等
不管用什么消息队列,第一道防线永远是业务幂等。什么是幂等?就是同一个操作执行一次和执行十次,最终结果是一样的。
我举一个最常见的例子:订单支付成功发消息,消费者收到消息后更新订单状态为“已支付”。如果不做幂等,重复消费时,第二次更新可能覆盖掉第一次之后的新状态(比如“已退款”被覆盖成“已支付”),这就是严重的 bug。
常用的幂等方案有三种:
- 数据库唯一约束:将消息中的业务主键(比如订单号、流水号)设为唯一索引,重复插入直接报错,捕捉后当作成功处理。
- 状态机校验:更新数据前先检查当前状态,只有符合前置状态的才允许更新。
- 去重表:维护一张消费记录表,记录每条消息 ID 与业务主键的映射,消费前先查重。
在实际生产里,我比较推荐“唯一约束 + 状态机”的组合。光靠消息 ID 去重有时候会有问题,因为同一业务操作可能会被不同消息触发,单纯用消息 ID 判断可能误伤。
6.3 Pulsar 提供的消费保障机制:ackTimeout、Nack 与重试
虽然业务幂等是根本,但 Pulsar 依然给了你不少工具去减少重复发生的频率。
合理配置 ackTimeout
java复制Consumer<String> consumer = client.newConsumer()
.topic("persistent://public/default/order-topic")
.subscriptionName("order-service")
.subscriptionType(SubscriptionType.Shared)
// 60秒内没有 ack,Broker 会重新投递
.ackTimeout(60, TimeUnit.SECONDS)
.subscribe();
这个参数的设置原则是:必须大于你业务处理时长的 P99,否则就会频繁触发误投。如果你发现线上重复变多,先别慌,去看看是不是 ackTimeout 设置得太小了。
用 Nack 代替 ackTimeout 处理临时失败
有些时候你希望给消费者更长的处理时间,但又不想禁用超时机制。这时可以用 Nack(Negative Acknowledge)显式告诉 Broker 这条消息我暂时处理不了,你再给我一次机会:
java复制consumer.negativeAcknowledge(message);
negativeAcknowledge 会立即触发重新投递,但要注意它默认会有重试延迟(默认 1 秒)。你可以通过 negativeAckRedeliveryDelay 设置间隔:
java复制client.newConsumer()
.negativeAckRedeliveryDelay(5, TimeUnit.SECONDS)
.subscribe();
把 ackTimeout 和 Nack 配合使用的姿势是:保留一个较大的 ackTimeout(兜底),业务处理失败时主动 Nack(而不是等超时)。这样既能快速重试,又降低了超时误判的风险。
开启重试 Topic 与死信 Topic
对于“重试几次还是失败”的消息,最好让它进入重试队列,达到最大重试次数后转入死信队列,而不是无限循环。Pulsar 对这块的支持很完善:
java复制Consumer<String> consumer = client.newConsumer()
.topic("persistent://public/default/order-topic")
.subscriptionName("order-service")
.subscriptionType(SubscriptionType.Shared)
.enableRetry(true)
.deadLetterPolicy(DeadLetterPolicy.builder()
.maxRedeliverCount(3)
.retryLetterTopic("persistent://public/default/order-topic-retry")
.deadLetterTopic("persistent://public/default/order-topic-dlq")
.build())
.subscribe();
配置之后,处理失败的消息会先进入重试 topic,Pulsar 会在指定延迟后投递回来;超过最大重试次数后,消息会进入 DLQ。这种机制下,消费逻辑不需要自己写重试循环,Broker 帮你管。
实操心得:我在项目里把 DLQ 设计成了“报警器”。专门有个服务监听所有 DLQ,只要里有消息进来,就意味着有业务一直处理失败,直接触发告警,人工介入排查。这比在业务代码里打日志高效得多。
6.4 批量接收与逐条确认的最佳实践
当你用 batchReceive 批量拉取消息时,最容易犯的错误是批量 ack:
java复制List<Message<String>> messages = consumer.batchReceive();
// 处理...
consumer.acknowledge(messages); // 错误的做法
acknowledge(List) 会一次性确认整个批次。如果这一批里有一条消息处理失败了,但你已经确认了整个批次,这条消息就悄悄丢了。正确做法是逐条确认:
java复制for (Message<String> message : messages) {
try {
process(message);
consumer.acknowledge(message);
} catch (Exception e) {
consumer.negativeAcknowledge(message);
}
}
虽然逐条确认会多一点网络开销,但在共享订阅模式下,这批消息大概率归属不同消费者,逐条确认是唯一不会误伤的方式。
7. 生产环境迁 practical Pulsar:我看重的几个运维与调优经验
代码层面聊完,最后补一些偏运维方向的实操经验。消息队列这种东西,运行一天两天看不出问题,跑上一个月,各种细节就全出来了。
7.1 管理 Topic、订阅与积压消息的命令速查
Pulsar 的命令行工具 pulsar-admin 一定要熟练。我日常用得最频的几个场景:
查看订阅列表和积压情况:
bash复制bin/pulsar-admin topics stats persistent://public/default/order-topic
这个命令的输出里,重点看 msgBacklog(积压消息数)和 blockedSubscriptionOnUnackedMsgs 这些字段。积压持续上涨,说明消费端跟不上了,该扩容或排查瓶颈了。
清理某个订阅(谨慎操作):
bash复制bin/pulsar-admin topics unsubscribe \
--subscription order-service \
persistent://public/default/order-topic
删除订阅后,该订阅的游标会消失,未消费的消息不会再有该订阅去消费,做这个操作之前一定要确认没有消费者在运行。
7.2 监控告警指标
Pulsar 的 Broker 通过 Prometheus 暴露指标,下面这几个指标我建议优先盯起来:
| 指标名 | 含义 | 建议告警阈值 |
|---|---|---|
pulsar_broker_msg_backlog |
积压消息数 | 持续大于预设阈值 |
pulsar_broker_storage_size |
存储用量 | 接近 BookKeeper 磁盘容量 80% |
pulsar_broker_subscription_blocked_on_unacked_messages |
因未确认消息过多而被阻塞的订阅 | 出现即告警 |
bookkeeper_server_ADD_ENTRY_REQUEST |
BookKeeper 写请求数 | 持续高位时关注磁盘 IO |
7.3 关于 Broker 参数的一个提醒
单机演示和真集群的参数配置完全不是一回事。生产环境里,我建议你至少关注这几个配置:
maxUnackedMessagesPerConsumer:单个消费者未确认消息上限,默认 50000。调太大会导致内存暴涨,调太小会导致吞吐上不去。maxUnackedMessagesPerSubscription:订阅整体未确认消息上限。超过之后 Broker 会阻塞向该订阅投递新消息。ackTimeout的 Redelivery 延迟:可以通过broker.conf调整重投间隔。
这些参数在低版本和高版本里位置可能有变化,实际操作时以你对应版本官方文档为准。
8. 结语与个人建议
写到这儿,Pulsar 的核心概念、快速上手、重复消费应对和运维要点都过了一遍。最后聊几句个人的选型心得。
如果你所在团队基础设施比较传统、网内环境封闭、对云原生没有强烈需求,Kafka 仍然是可靠的选择。但如果你已经在用 Kubernetes、追求弹性扩缩容、希望消息中间件能够提供更灵活的订阅模型和跨地域复制能力,Pulsar 的性价比会高很多。
我个人的体会是,Pulsar 的上手门槛并不比 Kafka 高多少,真正需要花时间的是理解它“存储与计算分离”带来的思维模式转变。一旦你接受这个设定,很多问题(分区扩容难、存储耦合、跨地域复制复杂)就不再是问题了。
如果你准备在生产环境上 Pulsar,我的建议是先不要追求大而全的集群架构,从单机或者三节点集群开始,把消息确认机制吃透,把幂等方案落地,再逐步扩大使用范围。这套路径,我们团队实际走了一遍,稳定性和收益都是超出预期的。
