
算算日子COSCon‘25 同场活动 Pulsar Developer Day 就剩三天了。今年几个技术群里早就在传议程消息中间件方向的专场活动能做到这种热度确实不多见。借着这个由头我打算把 Pulsar 这些年的一些核心设计、实践心得以及大家在群里最常问的几个问题比如“messageId 为什么长得那么奇怪”一次性捋清楚。不管你是准备去现场还是打算云围观这篇内容应该都能让你对消息中间件这件事有个更立体的认识。1. 先给还没入坑的朋友消息中间件为什么是必修课1.1 从单体应用到分布式消息队列解决的三件事如果你刚接触分布式系统可能还不太理解为什么大家都在聊消息中间件。我用最朴素的话来解释当你的系统从一个进程变成几十个服务互相调用时总会有一些请求是一瞬间涌进来的也有一些操作根本不需要让用户一直等着。消息中间件做的事本质上就是三件削峰填谷、服务解耦、异步处理。拿电商下单来举例。用户点下“提交订单”那一刻后端要处理库存扣减、优惠券核销、积分变动、发送短信通知、生成物流单……如果所有这些逻辑都同步跑完再返回“下单成功”高峰期接口耗时可能直接飙到几秒甚至超时。实际生产环境里大家会把短信通知、积分累计这类非核心链路扔进消息队列订单接口只管落库和发消息剩下的事情由下游服务异步消费处理。这既保证了用户体验也避免了大促时数据库被瞬间打垮。从单体到微服务的演进过程中消息队列几乎成了标配。但“标配”不等于“随便选一个就行”选型和架构设计直接决定了你未来三到五年的运维体感。1.2 选型号之前先搞清楚你的场景是哪种目前主流消息中间件各有一批忠实用户RabbitMQ 以灵活的路由和轻量部署著称适合内部系统间的事件通知Kafka 凭借超高吞吐和成熟的生态成为日志采集、数据管道的事实标准RocketMQ 在电商和金融场景里口碑不错事务消息做得很成熟而 Apache Pulsar这几年凭借“存储计算分离”和“多租户”两大王牌在云原生环境里增长非常快。选型时我一般会问自己四个问题吞吐量要求是多少量级是否要求消息可以回放、重跑是否需要跨地域复制扩缩容时能不能做到不停机如果团队规模小、业务简单RabbitMQ 或者直接用云厂商的托管队列就够了如果每天要处理几十亿条数据Kafka 或 Pulsar 才是真正的选项。Pulsar 的差异化优势在于它的 Broker 不存储数据只负责读写调度扩容时不需要搬数据而 Kafka 的 Broker 既管读写又管存储扩容往往伴随着 rebalance 和数据迁移痛点比较明显。2. Pulsar 这几年凭什么站稳脚跟2.1 存储计算分离把“分层”发挥到极致Pulsar 最核心的设计就是存储和计算分离。它底层依赖 Apache BookKeeper 作为分布式日志存储服务Broker 层不落盘数据所有的消息数据都交给 BookKeeper 集群管理。打个比方Kafka 像一家“前店后厂”的餐馆做菜和上菜是同一批人翻台率受限于后厨Pulsar 则像中央厨房模式前厅只管接单上菜菜品由统一的中央厨房配送哪个分店忙了多开几家前厅就行不需要重新备菜。存储计算分离带来的第一个红利是扩容优雅。Kafka 集群如果磁盘空间吃紧通常需要新增节点然后把 partition 副本迁移过去期间可能影响线上流量。Pulsar 只需要给 BookKeeper 集群加节点数据会自动做 rebalanceBroker 层毫不知情。第二个红利是读写分离更彻底。生产者和消费者的流量可以单独调度Broker 节点只维持连接和计算任务状态全部放在 BookKeeper。某个 Broker 挂了客户端会自动重连到其他 Brokersession 恢复成本极低。第三个红利是存储成本下降。Pulsar 从早期版本就支持分层存储Tiered Storage老数据可以自动从 BookKeeper 卸载到 S3 或 HDFS 这类廉价对象存储。Kafka 在这块的能力也在补强但 Pulsar 生来就是这么设计的配合上更顺滑。2.2 多租户和云原生基因不是噱头很多中间件在传统 IDC 时代根本不需要考虑多租户每个部门一套集群就完事了。但到了云原生阶段资源要共享成本要分摊权限要隔离多租户就成了刚需。Pulsar 的Tenant → Namespace → Topic三级模型天生就是为了共享集群而设计的。不同的部门可以共享同一套 Pulsar 集群但数据完全隔离权限可以细粒度控制配额可以分别管理。我记得有次分享会上一位维护过万级 Topic 集群的哥们说他最喜欢 Pulsar 的一点是“Topic 是轻量级的”。Kafka 里 Topic 多了以后分区数膨胀会拖垮 Broker 的整体性能Pulsar 因为存储和计算分离单集群承载百万级 Topic 都不是奇怪事这在做多业务接入时特别香。你不需要每接入一个新业务就申请一个新集群开个 Namespace 就够了。2.3 和 Kafka 的对比没有银弹只有适不适合我在这几年实际落地中两种中间件都深度用过。Kafka 在日志管道、大数据生态集成方面依然无可替代计算引擎和 Kafka 的集成深度远超 Pulsar。如果你整个技术栈都围绕 Flink、Spark、ClickHouse 转Kafka 可能是心智负担最低的选择。但如果你在乎的是云原生部署、跨地域容灾、多租户隔离、以及未来可能出现的突发扩容需求Pulsar 的架构优势就会逐渐体现出来。尤其是跨地域复制Pulsar 原生支持多集群的异步复制配置而 Kafka 的 MirrorMaker 一直有种“外挂工具”的糙劲儿。正如没有银弹一样选型时列个表格把你的业务场景、团队熟悉度、未来规划填进去该选谁答案自然浮现。3. 热词解析Pulsar 的 MessageId 为什么会是这个样子3.1 先来破解那一串字符的秘密兄弟群里有人甩了个问题“为什么我拿到一个 messageId 长这样messageId|28077:20854:-1这到底是个啥” 我第一次看到这个格式的时候也愣了一下因为习惯 Kafka 的 offset 是单调递增的整数一眼能看懂。Pulsar 的 MessageId 却是一串“ledgerId:entryId:partitionIndex”的结构。要理解这个格式得从 Pulsar 的存储结构说起。我在上面提到Pulsar 用 BookKeeper 来存数据而 BookKeeper 最基础的存储单元是Ledger。Ledger 是一段追加写入的日志文件里面包含一条条 Entry。对 Pulsar 来说一个 Topic 的消息会顺序写入一系列 Ledger 中每条消息在 Ledger 里对应一个 Entry。所以28077是Ledger ID消息属于哪一个 Ledger20854是Entry ID消息在 Ledger 里的物理偏移序号-1或0是Partition Index表示这条消息来自哪个分区-1 通常用于非分区 Topic 或管理场景。所以messageId|28077:20854:-1的意思是这条消息位于第 28077 个 Ledger 的第 20854 条 Entry 上这条消息所在的 Topic 是一个非分区主题。用文件系统来类比的话Ledger 就像一本书Entry 就是书里的页码两者定位才能精确找到一条消息。3.2 Ledger 机制读懂了它就读懂了 Pulsar 一半Ledger 是 Pulsar 存储的一个核心抽象。每个 Ledger 都有几个特性追加写入、不可变性、按 Entry 序号随机读取。一个 Ledger 写满一定大小或者存活超过一定时间后就会自动关闭然后新建一个 Ledger 继续写。这个过程叫 Ledger Rollover。为什么要设计成小段小段的 Ledger而不是一个大文件写到底这里面的门道很深。第一个好处是便于容量管理和恢复如果某个 Bookie 节点挂了只需要对它负责的那部分 Ledger 做数据恢复而不是整库重建。第二个好处是便于实现高效的 TTL 删除消息过期后整段 Ledger 直接删掉就行不需要像 Kafka 那样做日志紧凑和分段清理。第三个好处是并行度更好多个 Ledger 可以分布在不同的 Bookie 上写入流程天然可以并行扩展。这里还要提一个隐藏机制——Ensemble。Pulsar 写入时不是把所有副本都写在同一批 Bookie 上而是会为每个 Ledger 动态选一组 BookieEnsemble数据以条带方式分布写入。这个设计保证了当某个 Bookie 出问题时数据依然是可读的而且不影响其他 Ledger 的写入。生产环境里最常见的配置是 E3, W3, A2意思是 3 个 Bookie 组成一个 Ensemble需要写入全部 3 份才算成功但允许 1 个节点故障不影响读取。这套逻辑源自 BookKeeper 的 Quorum 机制理解了你就能解释为什么 Pulsar 能在高可用和数据一致性之间维持不错的平衡。3.3 为什么不能直接用一个自增整数当 offset很多人会问Kafka 的 offset 是分区内单调递增的整数直观又简单Pulsar 为什么要搞这么复杂这个问题的答案还是要回到存储计算分离上。Kafka 的 offset 是分段日志里的位置索引它跟本地磁盘上的文件偏移强相关所以天然只能由某个 Broker 自己管理消费者恢复位点时必须找对 Broker。而 Pulsar 的消息位点由 LedgerId EntryId 唯一定位Broker 层不存储实际数据消费者无论连接哪个 Broker都能通过这个 ID 去 BookKeeper 里读取。这种设计让 Pulsar 的客户端连接可以无缝漂移Broker 故障时消费者几乎无感知。此外Pulsar 的 MessageId 还包含 Batch 概念。生产端开启批量后多条消息会打包进同一个 Entry但对外暴露的 MessageId 依然能区分出 Batch 内部的单条消息。比如(ledgerId, entryId, partitionIndex, batchIndex)四个维度比 Kafka 的“offset batch 内相对位置”要更清晰。刚接触的人可能会觉得格式复杂但当你开始做消息回溯、精确消费到某一条消息时这种设计带来的准确性是简单整数 offset 做不到的。3.4 从 MessageId 衍生出去Cursur 与消息确认机制理解了 MessageId 之后“游标Cursor”的概念就比较好懂了。Pulsar 的每个订阅都有一个游标游标里存储的是这个消费者组当前消费到的 MessageId。消费成功后客户端会发送 ACK 给 BrokerBroker 把游标前移。这个游标信息默认保存在 BookKeeper 里叫作Cursor Ledger相当于把消费进度也做成了高可用存储。实际排查问题的时候游标和 MessageId 的配合非常有用。比如消费者组堆积了你可以用pulsar-admin topics stats命令查看msgBacklog那是当前游标到最大 MessageId 之间的消息条数。也可以用peek-messages命令指定 MessageId 来查看某条历史消息的内容这在定位“某条消息到底有没有被消费到”的场景下特别好用。这些能力在 Kafka 里实现起来很别扭因为消息位点放在消费者端服务端对堆积状态的管理要弱不少。4. 生产环境绕不开的订阅模式、分层存储与运维要点4.1 三种订阅模式怎么选Pulsar 的订阅模型是它的又一大卖点。它原生支持三种消费模式很多人刚接触时容易混淆我在这里用一个生活场景来拆解独占订阅Exclusive一个 Topic 同一时刻只能有一个消费者适合强顺序场景比如把订单状态流转消息按顺序处理不能并发。这种模式最简单但扩展性最差。共享订阅Shared多个消费者共同消费一个 Topic 的消息消息按 round-robin 或 pending-ack 状态分配给不同消费者吞吐量最高但消息顺序性无法保证。适合大多数数据处理场景比如消息量很大、处理逻辑之间没有严格依赖关系。灾备订阅Failover多个消费者中只有一个活跃消费者接收消息其他消费者作为备用活跃消费者挂掉后自动切换。相当于带高可用的独占模式。还有一种是 Key_Shared 订阅介于 Shared 和 Exclusive 之间相同 key 的消息只发给同一个消费者比如把同一个用户 ID 的所有订单事件固定发给一个消费者处理这样既能并行消费又能保证单用户的局部有序。生产环境里我用 Key_Shared 解决过一个典型问题用户积分变动必须按时间顺序处理但不同用户之间可以并行用 Exclusive 太浪费用 Shared 会乱序Key_Shared 刚好完美解决。4.2 分层存储省钱大法Pulsar 的分层存储Tiered Storage是个非常实用的功能但宣传得还不够。默认情况下数据只存在 BookKeeper 里而 BookKeeper 一般用的是 SAS 盘或者 SSD。数据量大了以后存储成本非常扎眼。Pulsar 可以配置自动卸载策略把超过一定时间或者一定容量的旧数据搬到 S3、阿里云 OSS、腾讯云 COS 这类对象存储里读的时候如果 BookKeeper 里没有会自动从对象存储拉回来。这个机制在“消息回溯”场景里特别有用。很多团队要求消息至少保留 7 天甚至 30 天以便排查问题或者重新跑数。如果不做分层存储你就得为这 30 天的数据预留大量高性能磁盘。开了分层存储之后热数据在 BookKeeper冷数据在对象存储成本能降一个数量级。我见过有团队把 Pulsar 保留了 90 天的数据存储成本跟 Kafka 保留 3 天差不多这就是分层存储的魔力。4.3 运维踩坑几个值得记录的真实案例Pulsar 的整体运维体验比传统消息中间件要好但也不是没有坑。我把自己落地过程中遇到比较多的问题整理成一张速查表现象可能原因排查方法生产端发送延迟突增BookKeeper 写入延迟变高磁盘 I/O 忙bookkeeper shell ledgercheck检查磁盘状态观察 bookie 的 Journal 和 EntryLog 写入延迟消费端 backlog 持续增加消费者处理慢某个消费者挂掉没有重连用pulsar-admin topics stats看 backlog检查客户端日志连接状态Topic 写入失败 “No such ledger”Ledger 被自动删除但客户端还在写时间窗口问题调整保留策略确保写入过程中 Ledger 不会被清理客户端连接不断重连Broker 和 Bookie 之间的网络抖动认证过期检查认证 token 有效期查看 broker 日志里的连接错误分区数量无法减少设计阶段分区数不合理先评估消息量和消费并发度一般先小后大确实不够再扩容分区大家最容易踩的坑其实是“分区数设太多”。Pulsar 对 topic 数量的容忍度极高但分区太多会导致 BookKeeper 的 Ledger 数量爆炸增加 ZooKeeper 和 Bookie 的元数据压力。我一般建议初期按消费者并行度来定分区数后期根据实际流量再逐个加不要未雨绸缪地建一堆分区。4.4 客户端参数调优三个我建议必调的配置Pulsar 客户端很多参数都有默认值但生产环境里默认值往往不是最优解。第一个是生产端的 Batching。默认 batching 是开启的但如果你追求极低延迟比如控制在 5ms 以内需要把batchingMaxPublishDelay调小或直接关闭 Batching。反过来如果要追求高吞吐可以把batchingMaxMessages调到 1000 以上batchingMaxBytes调到 128KB 以上。这个取舍跟 TCP 的 Nagle 算法有点像具体按你的业务对延迟的敏感度来。第二个是消费端的 Ack 超时。很多人设了ackTimeout但设得太短比如 1 秒消费逻辑稍微慢一点就会触发重投导致消息重复处理的概率上升。如果你下游不要求精确一次建议把 ackTimeout 调大甚至设为 0不超时配合消息重试队列去处理失败场景。注意ackTimeout 和negativeAckRedeliveryDelay是两套独立的重试机制别混淆。第三个是接收队列大小。receiverQueueSize默认 1000对处理很快的逻辑来说这个值偏小容易让消费者频繁陷入等待。我通常在 CPU 密集型处理时设为 100~200IO 密集型或调用下游接口场景设到 1000 或更高避免消费者线程空转。5. 活动背面的技术信号为什么开发者日值得你专门跑一趟5.1 从议程看行业趋势说实话国内专门围绕 Pulsar 做线下开发者日的机会并不多一年到头可能也就几场。去年我在一次类似的活动中听到了好几个非常有启发的议题有团队分享了怎么把 Pulsar 用在金融级交易系统里保证消息不丢不重也有人讲了怎么用 Pulsar Functions 做轻量级流处理替代一部分 Flink 任务还有人在分享 Pulsar Lakehouse 的实践把消息中间件和数据湖打通的路径已经非常清晰。从这些议题能明显感觉到Pulsar 已经不是早期那个“技术超前但生态不足”的项目了。它开始向金融、运营商、车联网这些传统行业渗透在云原生、数据集成方面也出现了一大批成熟的同场产品。今年 COSCon’25 同场活动定名为 Pulsar Developer Day来的 speaker 和 sponsor 大多是实际落地的一线工程师不是空谈架构的纸上谈兵这种内容密度比主流技术大会的高层 keynote 要实在得多。5.2 带着问题去收获才更大参加这类活动的正确姿势不是坐在台下听 PPT而是带着自己的实际问题去。我每次去都会准备一张问题清单比如我们在生产环境遇到 Bookie 扩容后写入抖动官方有没有推荐的 rebalance 策略在 Pulsar 里实现“延迟消息”和“定时消息”最佳实践是什么Pulsar 和 Flink 集成的 checkpoint 一致性如何处理在跨地域复制场景下消息延迟的监控指标怎么设阈值这些问题在现场很容易通过 speaker 的分享和 QA 环节找到答案甚至可以当面加微信交流后续细节。线上文档写得再多也不如和核心维护者、一线实践者面对面聊十分钟收获大。今年活动倒计时只剩三天票务信息在 COSCon 官网可以看到对消息中间件生态感兴趣的朋友我觉得完全可以抽出半天时间过去转转。5.3 如果你没法到场可以这样保持同步对于没法到场的朋友我的建议是重点留意活动后的 PPT 和视频回放一般 COSCon 系列的产出质量高官方渠道都会公开。同时把活动中提到的实践项目、仓库链接都收藏起来动手跑一遍比收藏一堆“技术文章”有用得多。另一个做法是加入 Pulsar 中文社区或邮件列表很多讨论从活动当天会一直延续到线上你在群里提问经常能得到活跃贡献者的直接回复。6. 写在最后我的个人体会我自己从 Kafka 转向深入使用 Pulsar大概经历了一个从“这玩意怎么这么复杂”到“原来这样设计是合理的”再到“回不去了”的过程。MessageId 的复杂结构、Ledger 的切分机制、游标的持久化初看都是增加理解成本的东西但用久了会发现这些设计全都是为了支撑同一个目标让存储和计算解耦让扩缩容不再痛苦让大数据量下的消息系统保持稳定。马上就是 Pulsar Developer Day 的日子了这种主题的交流机会且行且珍惜。不管你是消息中间件的老兵还是刚刚入行的新人去听听一线踩坑的人怎么说远比自己在电脑前瞎琢磨效率高。准备去现场的朋友我建议提前列出你当前系统里三个最痛的问题直奔对应议题的讲师别不好意思问。三天后见。