news 2026/9/29 16:35:16

Pulsar核心机制与生产实践:从选型到排障一次讲透

作者头像

张小明

前端开发工程师

1.2k 24
文章封面图
Pulsar核心机制与生产实践:从选型到排障一次讲透

COSCon‘25 同场活动的 Pulsar Developer Day 议程正式发布之后,我朋友圈里几个做中间件运维的朋友几乎同时转了同一条消息。消息中间件这个领域,讨论热度一直没有降过,尤其是 Pulsar 这种计算存储分离风格的系统,在选型会上总被拿来和 Kafka 对比。作为一个从 Pulsar 2.x 版本就开始在业务线里实际使用的人,这篇我不打算复述官方议程有多少个议题,而是借着这份 Pulsar Developer Day 议程里大家普遍关心的方向,把消息中间件选型、Pulsar 核心机制、上手路径、生产排障这四件事一次讲透。适合谁看?给团队做消息中间件选型的架构同学,以及已经在用 Pulsar 但想补一补内部原理和排障思路的研发与运维。

1. 为什么 Pulsar 在消息中间件选型里越来越被认真对待

1.1 存储与计算分离:Pulsar 和 Kafka 的根本差异

Pulsar 最核心的设计,是把 Broker 和 BookKeeper 拆成两层。Broker 是无状态的接入层,负责处理生产消费协议、管理 topic 元信息和各种策略;BookKeeper 是真正的存储层,负责持久化消息数据。Producer 把消息发到 Broker 后,Broker 将消息追加到 BookKeeper 的 ledger 里,而不是写进本地磁盘。这和 Kafka 把分区日志直接挂在 broker 本地磁盘是两种完全不同的存储模型。

只有单机或三五台机器的小集群,可能感受不到这个差异。等集群规模上来,存储计算分离的好处就非常明显了:扩 Broker 不需要搬数据,只需要把 namespace bundle 的归属重新划分;缩容也同理。Kafka 扩容分区往往要经历数据重分布,Pulsar 的 Broker 本身不持有数据,水平扩缩容更像“加收银台”,店里的仓库不动。

再展开一点,BookKeeper 里的数据可靠性由几个参数共同决定:ensemble 指定数据分布到几个存储节点,write quorum 指定同一份数据写几个副本,ack quorum 指定必须有几个副本确认才能返回成功。默认配置下是三个存储节点各持一份,写入时等其中两个确认就返回。这个机制保证了单个存储节点故障时数据仍然可读,也解释了为什么 Pulsar 的可用性模型和 Kafka 不完全一样。

1.2 三种订阅模式:不是概念而是消费模型

Pulsar 的订阅模型是它辨识度很高的地方。exclusive 是排他消费,一个订阅同时只允许一个消费者在线,消息严格有序,适合订单状态机、审计日志这类“必须按顺序处理”的场景;failover 是主备模式,主消费者挂了之后备消费者接管,消费顺序有约束但可用性更高;shared 模式下多个消费者竞争消费,吞吐高但顺序没有保证,适合通知推送、数据灌仓这些对顺序不敏感的场景。后来提供的 key_shared 让相同 key 的消息固定路由到同一个消费者,在并发和按 key 有序之间找到了平衡。

实际选型时,这三种订阅模型让你在一套系统里同时实现队列和流两种语义。shared 订阅就是消息队列,exclusive/failover 订阅就是有序流处理,这也是 Pulsar 常说的“one platform for queuing and streaming”的含义。团队里如果既有异步削峰场景,又有实时流处理场景,用一套中间件统一承接会省掉很多跨团队沟通和双栈运维成本。

1.3 对照 Kafka 与 RocketMQ:适合场景与选型边界

单纯聊架构可能不够直观,我把几个主流中间件的差异整理成一张表,信息来自自己实际使用和社区公开资料,大方向足够支撑选型判断。

维度PulsarKafkaRocketMQ
存储架构计算与存储分离(Broker + BookKeeper)存储耦合在 broker 本地分区日志存储耦合在 broker 本地 CommitLog
水平扩容扩 Broker 节点即可,存储不搬数据分区扩容需要数据重分布主要靠节点垂直扩容或重建主题
消费模型exclusive/shared/failover/key_shared 多模型消费组模式,偏流式消费Push/Pull 加队列模型,偏消息队列
多租户隔离tenancy/namespace 原生体系需要结合 ACL 和配额自建支持资源隔离,机制相对重
典型场景混合负载、多业务共享集群、跨地域复制大流量日志管道、事件流处理电商、订单、延时消息等传统队列场景

选型边界我个人的判断是:如果业务形态非常单一,链路就是“数据进 Kafka、实时引擎消费”,周边组件也多依赖 Kafka 生态,那继续用 Kafka 是稳妥的选择;如果需要在同一个集群里同时支撑高吞吐队列和有序流处理,又希望多个团队共享基础设施时资源隔离不互相干扰,Pulsar 的优势会体现得更充分。另外,国内很多团队已经有深入的 RocketMQ 使用沉淀,那类场景里迁移的摩擦成本也需要认真评估。

顺带说一句,消息中间件的范畴比很多人想象中大得多。像嵌入式飞控场景里的 uORB,它在控制回路里承担着传感器、控制模块之间高频数据搬家的职责,和 Pulsar 这种面向高吞吐、多租户的分布式系统不是一个重量级,但本质都是“数据的搬运方式”。所以选型时永远先想清楚场景约束:数据规模、延迟要求、团队运维能力分别是什么,再决定用哪种形态的中间件。

2. 消息进出都经历了什么:Pulsar 核心机制拆解

开发者日这类活动,最有价值的通常不是功能宣传,而是把内部机制摊开来讲的议题。下面这些点是我认为理解 Pulsar 时必须掌握的,也是排障时绕不开的底层逻辑。

2.1 从 Topic 到 Ledger:一条消息的落盘路径

Pulsar 的 topic 是逻辑队列,topic 下可以有多个 partition。每个 partition 在存储层不是挂在独立物理文件上,而是对应一条或多条 BookKeeper ledger。ledger 是 append-only 的日志结构,消息写入后不断增长,到了滚动条件会生成新的 ledger,旧的 ledger 只有在所有订阅游标都越过去、且保留策略允许时才会被回收。

一条消息的物理落盘流程大致是:Producer 发送到 Broker,Broker 把消息按批组织成 BookKeeper entry,追加写到底层 ledger。存储节点完成写入并确认,等 ackQuorum 条件满足后,Broker 才向 Producer 返回成功。看起来只是“发送-确认”两步,背后却牵涉多个存储节点的副本写入。

为什么 ledger 要做成分段滚动而不使用一个永远增长的大文件?因为分段让副本恢复变得更容易,某个节点坏掉后只需要定位和补齐缺失的那一段,数据回收也可以按段粒度处理,不用对超大文件做全局合并。理解这一点,再看 Pulsar 的存储模型就不会觉得它只有“计算存储分离”一个标签。

2.2 Cursor 与 Backlog:为什么消息堆在队列里不消失

每个订阅都有独立的 cursor,cursor 记录当前已确认消费到哪条消息。Broker 判断一条消息能不能删除,不只看 topic 的保留策略,还要看该 topic 下所有订阅的 cursor 是否已经越过它。很多人设置了 retention 却发现磁盘迟迟不释放,原因就在这里——某个订阅的游标还没推进到位,消息就不能从 ledger 中删掉。

Backlog 指的是某个订阅尚未确认的消息数量。这里有个常见的认知误区:以为 topic 存储很大就等于 backlog 很高。其实同一 topic 可以挂多个订阅,不同订阅的消费速度相互独立。A 订阅的消费者早就追平,B 订阅因为消费脚本挂掉而卡住了几千条,topic 的 StorageSize 仍然会保持高位。排障时建议按订阅维度看指标,用 admin 命令查看每个订阅的 backlog 总数,而不是盯着整体容量猜测。

2.3 留存策略与分层存储:存储成本控制的实践

Pulsar 支持分层存储,可以把已经确认过的消息转储到对象存储上,比如 S3、OSS,Broker 本地只保留热数据。这个能力对“消息要保留很久,但不需要低延迟访问”的场景特别有用,典型就是审计日志、离线回放、合规留档。数据卸载之后,热存储压力下降,成本曲线会好看很多。代价是读取时明显变慢,对象存储的延迟比本地热盘高,低延迟实时链路不要依赖它兜底。

保留策略配置上,retention 控制已确认消息的保留时长,TTL 控制未确认消息多久后被自动标记。二者配合一旦失误,生产事故来得很快。我见过一个团队把 retention 设成 -1,高吞吐 topic 的所有历史消息全部永久留在热存储里,磁盘成本直接失控。经验是:保留时长必须写成具体数值,并且定期复盘主题的存储量,发现持续增长要立刻查策略而不是继续堆机器。

2.4 交付语义与消费幂等:重复消息避不开

Pulsar 默认的投递语义是至少一次,消费者处理完消息后如果 ack 超时或网络异常,Broker 会重新投递。这意味着消费端一定会遇到重复消息,设计业务逻辑时必须把幂等考虑进去,而不是事后手忙脚乱地补数。

比较简单的方案是给每条业务消息带上全局唯一 ID,消费端在处理前先查是否已经处理过;复杂链路可以引入存储层唯一约束,比如订单号、流水号作为去重键。Pulsar 也提供了事务能力,支持多消息、多分区之间的原子写入与确认,可以实现更严格的恰好一次语义,但事务对客户端版本、服务端配置都有要求,使用门槛不低。大多数业务场景,先把幂等表做好,比直接上事务要更务实。

3. 5分钟跑通 Pulsar 并完成客户端接入

理论说再多都不如动手跑一遍。下面这套流程完全可以作为第一次接触 Pulsar 的上手路径,从零到能收发消息只需几分钟。

3.1 本地环境搭建:Docker 单机模式

我习惯用官方镜像把 Pulsar 以 standalone 模式跑起来,一条命令就能完成:

docker run -it -p 6650:6650 -p 8080:8080 apachepulsar/pulsar:3.3.0 bin/pulsar standalone

端口 6650 是客户端通信端口,8080 是 admin REST 接口。启动日志里看到“messaging service is ready”之类的输出就说明服务已经就绪。第一次跑的时候最大的坑是只映射了 8080 端口,客户端连不上,盯着防火墙排查半天才发现端口少了一个。如果本机没有 Docker,也可以下载官方二进制包,解压后直接执行 bin/pulsar standalone,效果一样。网络慢的同学建议提前把镜像拉好。

standalone 模式只适合本地验证和开发联调。生产环境至少要部署三个 broker 和三个 bookie,并单独配置元数据服务,不能用单机模式扛业务流量,毕竟它没有故障转移能力。

3.2 用命令行完成一次完整收发

服务起来之后,先用 admin 命令确认集群状态,再创建 topic,然后生产和消费:

bin/pulsar-admin clusters list bin/pulsar-admin topics create persistent://public/default/my-topic bin/pulsar-client produce my-topic --messages "hello-pulsar" -n 10 bin/pulsar-client consume my-topic -n 10

bin/pulsar-admin clusters list 会返回当前集群的元数据服务列表;topics create 命令在 persistent://public/default/my-topic 下创建一个持久化主题。produce 命令会连续发送 10 条消息,每条内容就是 hello-pulsar;consume 命令启动后会拉取当前订阅里的消息,消费完 10 条后进程仍然挂着等待新消息,按 Ctrl+C 退出即可。命令行跑通,说明环境、端口、认证都没有问题,才有必要进入客户端代码阶段。如果某个端口被占用,或者容器内配置文件没挂出来,命令会直接提示连接失败,这类问题优先检查端口映射而不是代码。

3.3 客户端代码:从发送到确认的关键写法

Java 最小可运行示例:

import org.apache.pulsar.client.api.*; PulsarClient client = PulsarClient.builder() .serviceUrl("pulsar://localhost:6650") .build(); Producer<String> producer = client.newProducer(Schema.STRING) .topic("my-topic") .enableBatching(true) .batchingMaxMessages(1000) .batchingMaxPublishDelay(10, TimeUnit.MILLISECONDS) .create(); producer.sendAsync("hello-pulsar".getBytes()); Consumer<String> consumer = client.newConsumer(Schema.STRING) .topic("my-topic") .subscriptionName("my-sub") .subscribe(); Message<String> msg = consumer.receive(); try { System.out.println(msg.getValue()); } finally { consumer.acknowledge(msg); } client.close();

Python 版本更短:

from pulsar import Client client = Client('pulsar://localhost:6650') producer = client.create_producer('my-topic') for i in range(10): producer.send(('msg-%d' % i).encode('utf-8')) consumer = client.subscribe('my-topic', 'my-sub') msg = consumer.receive() print(msg.data().decode('utf-8')) consumer.acknowledge(msg) client.close()

客户端代码里最容易出问题的不是连接本身,而是 ack 行为。生产者发送建议用 sendAsync 而不是同步 send,同步发送在高峰期会阻塞调用线程。消费者 receive 到消息后必须在 finally 里 acknowledge,尤其是异常分支也要保证 ack 或 nack,否则消息会不断重投。shared 订阅下未确认消息超时后 Broker 会把它们重新分发给其他消费者,业务侧看到的重复消息比例会显著上升。

3.4 几个关键配置参数:从服务端到客户端

Pulsar 的配置项很多,我不打算罗列流水账,挑几个直接影响线上表现的说说。

maxUnackedMessagesPerConsumer 决定单个消费者最多可以有多少条消息处于未确认状态。shared 订阅下如果消费者一直不 ack,超过这个阈值后 Broker 会停止继续推送。调大这个值可以提升单个消费者的处理空间,但也意味着有更多消息没有及时确认,消息丢失风险窗口会变大。

sendTimeout 是生产者发送超时时间,默认 30 秒。网络抖动时频繁触发超时,先确认超时阈值设置是否合理,再看发送方到 Broker 的链路质量,而不是一上来就把 timeout 调到几分钟。同步 send 和异步 sendAsync 对超时表现也不同,后者更能抗住瞬时抖动。

enableBatching 控制是否批量发送,默认开启,适合高吞吐场景。batchingMaxMessages 和 batchingMaxPublishDelay 需要同时权衡,批处理窗口太长,尾部延迟会明显拉高。给业务设定目标延迟,再倒推批参数,比盲目调大要靠谱。

服务端还有个容易被忽略的配置是 defaultNumberOfNamespaceBundles,它决定 namespace 下 bundle 的数量。bundle 数量充足,topic 在 Broker 之间的分布会比较均匀,不容易出现某个 Broker 上热点累积;数量太大,元数据同步和查找压力又会上升。不同集群之间很难直接抄作业,先根据 topic 总数和 broker 数取一个保守值,再观察负载分布逐步微调。

4. 生产环境 Pulsar 常见问题与排查技巧实录

这一章写的是我在实际维护中见过、踩过的典型问题。Pulsar 社区现在活跃度不低,遇到奇怪问题换个搜索姿势也常能在社区 issue 里找到蛛丝马迹,但更关键的还是先有一套自己的排查节奏。

4.1 问题速查表:现象、原因、处理建议

现象可能原因处理建议
Producer 发送超时topic 自动创建被关闭、Broker 到 BookKeeper 连接异常、磁盘容量不足检查 allowAutoTopicCreation 配置,看 Broker 日志与存储节点状态,确认磁盘水位
Backlog 持续增长但消费者在线消费逻辑 ack 缺失、单条消息处理耗时过长、反序列化异常统计消息处理耗时,检查 receive 后是否调用 ack,看消费者线程池情况
消费者频繁被断开心跳超时、客户端与服务端版本行为差异调大心跳间隔,统一客户端版本,查看 Broker 端连接日志
topic 磁盘占用异常增长retention 或 TTL 配置不当、某个订阅消费停滞查看 StorageSize 曲线,按 namespace 收敛策略,检查各订阅 cursor 位置

排查时我习惯从下往上拆:先确认 topic、namespace、broker 三级指标没有异常,再看客户端日志,最后检查硬件和网络层。很多表面上是中间件故障的问题,最后查出来都是消费端代码的锅。

4.2 一次 Backlog 暴增事故:从现象到根因的完整排查

挑一个印象最深的场景讲讲。某个下午业务方反馈通知延迟严重,后台数据显示某 topic 的 backlog 从几千条一路涨到几百万条。我第一反应不是去重启消费者,而是先拉生产端和消费端指标。生产速率没有异常,但消费组 Pod 的内存持续走高,反复触发 OOM,重启后的消费者要从上游 checkpoint 重新拉取数据,又赶上业务高峰期,backlog 开始滚雪球。

处理过程分了三步。先扩消费端副本数到平时的两倍,同时关掉非核心日志打印降低内存压力;等 backlog 追平后再逐步缩容到正常水位。然后排查 OOM 根因,发现是消息体比预估大很多,反序列化时缓存了过多对象,GC 跟不上。最后把消费端内存上限调整,并给这个 topic 单独加了一条 backlog 监控告警。

这次排障给我最大的启发是:backlog 暴增只是一种现象,背后的根因完全可能不在中间件本身。如果你一开始就把注意力放在调配置、加机器上,很容易掩盖真正的问题,等下一个业务高峰再复发。

4.3 几个文档里不会写的团队踩坑心得

Namespace 规划一定要从第一天就认真做。建议按“团队/环境/业务域”拆分,比如 order-prod、warehouse-test,权限、配额、限流都可以挂在 namespace 上,后续多租户隔离会非常省事。等 topic 数量上百之后再重新拆 Namespace,涉及大量迁移和客户端改动,比一开始就规划好要痛苦太多。

批量参数不要贪。batchingMaxMessages 和 batchingMaxPublishDelay 不是越大越好。之前一个业务团队为追求吞吐把批消息数调到十万,偶发批量发送延迟飙升,排查半天发现批处理窗口过长,几条慢消息拖住了整批。根据目标延迟反推配置值,比盲目拉大参数更靠谱。

升级客户端版本前先做兼容性测试。Pulsar 服务端小版本升级通常平滑,但大版本之间有些默认行为和指标口径会变,预发环境完整跑一遍消费流程,比上线之后再回滚省心得多。有一次我们大版本升级后老客户端的重连频率明显上升,回退客户端版本才恢复正常,从那以后我就把版本兼容性测试写进了发布规范。

监控不要只建“有没有”的指标。按 namespace、topic、订阅三个维度分别看,backlog 曲线、StorageSize、生产消费速率这几个基础指标全都要有。之前有一次集群总体指标很健康,只有一个核心订阅的 backlog 在悄悄上涨,基础监控完全看不出来,最后业务方投诉了才发现。三层维度的监控面板,定位问题的效率完全不一样。

如果你也准备在团队里推进 Pulsar,我个人的建议是从一个小规模业务开始,先把 namespace 结构、监控面板、保留策略这些基本功做扎实,再逐步放大。等哪天社区开 Pulsar Developer Day 这类活动,你会发现最值回票价的反而不是新功能,而是那些真正在生产环境趟过坑的人讲的使用细节。消息中间件这种底层设施,可靠性和可维护性比炫技更值钱。

版权声明: 本文来自互联网用户投稿,该文观点仅代表作者本人,不代表本站立场。本站仅提供信息存储空间服务,不拥有所有权,不承担相关法律责任。如若内容造成侵权/违法违规/事实不符,请联系邮箱:809451989@qq.com进行投诉反馈,一经查实,立即删除!
网站建设 2026/9/29 16:35:06

Starnet概念解析:技术命名规范与术语验证方法

我无法根据当前输入生成符合要求的博文。原因如下&#xff1a;项目标题为“starnet”&#xff0c;但项目正文、关键词、摘要描述全部为空&#xff1b;所谓“相关热搜词”和“最新网络热词”仅重复出现“starnet”一词&#xff0c;无任何上下文、定义、领域指向或可验证信息&…

作者头像 李华
网站建设 2026/9/29 16:34:59

PostgreSQL数据库大小查询指南:函数原理、实操SQL与DeepSeek排障实践

上周半夜两点&#xff0c;监控告警把我从床上叫醒&#xff1a;磁盘使用率92%。爬起来第一件事不是去看日志&#xff0c;而是登录数据库&#xff0c;想搞清楚到底是哪个库、哪张表把空间吃没了。这个场景做运维的应该都经历过&#xff0c;而PostgreSQL在这一点上确实非常友好——…

作者头像 李华
网站建设 2026/9/29 16:34:58

F28377D硬件加速器实战:TMU/VCU-II/CLB详解

做电机控制和数字电源这些年&#xff0c;F28377D这块片子我前前后后用过几个项目。说句实在话&#xff0c;很多人把它当普通C2000用&#xff0c;跑跑主频、调调PWM、做做ADC采样&#xff0c;核心的浮点运算靠CPU硬扛。但这颗芯片真正值钱的地方&#xff0c;其实是它在C28x内核旁…

作者头像 李华
网站建设 2026/9/29 16:33:54

推测执行详解:从Hadoop MapReduce到Spark的调优实战

1. 一次真实的集群“掉队”事故&#xff1a;我为什么开始重视推测执行大概两年前的这个时候&#xff0c;我负责的一个离线数仓集群出了个诡异现象&#xff1a;每晚跑核心ETL任务&#xff0c;整个DAG都跑完了&#xff0c;就卡在最后几个MapReduce job上。点开Hadoop Application…

作者头像 李华
网站建设 2026/9/29 16:33:50

AI工程化实战路线:从数据清洗到模型部署的完整指南

说实话&#xff0c;我第一次听到「AI工程师」这个称呼的时候&#xff0c;自己先在心里打了个问号——这不就是调模型的人吗&#xff1f;但真正扎进去做了几年之后才发现&#xff0c;ai-engineering 跟我们对 AI 的浪漫想象完全是两码事。它不是写几行代码、跑一个训练脚本那么简…

作者头像 李华
网站建设 2026/9/29 16:33:28

办公RAG系统实战:LangChain+FastAPI+Vue3生产级AI OA

1. 这不是又一个“大模型OA”的Demo&#xff0c;而是能真正跑通的智能办公最小可行系统我去年带三个本科生做毕业设计&#xff0c;其中两个组选了“AI办公系统”&#xff0c;结果第一周就卡在环境里&#xff1a;有人装了三天Vue3还跑不起来前端路由&#xff0c;有人用LangChain…

作者头像 李华