尧图网站设计 尧图网站设计YAOTU DESIGN
ARTICLE DETAIL

资讯详情

深耕网站设计与一线实操的经验洞察。

Pulsar开发者日复盘:存算分离、延迟消息与消息查看实践

Pulsar开发者日复盘:存算分离、延迟消息与消息查看实践 当COSCon25会场大屏打出“Make MQ Great Again”这个标题时台下一片笑声。作为一个常年和消息队列纠缠的开发者我太懂这个双关背后的情绪了MQ消息队列这个领域已经很久没有被这样“认真对待”过。过去十年Kafka几乎成了消息队列的代名词选型时一句话就能结束讨论。但Pulsar偏要站出来用存算分离、多租户、分层存储这些底子把“消息队列应该是什么样”重新摆在桌面上。这篇文章是这次Pulsar Developer Day 2025的回顾也是我在现场听完一天分享、被拉着讨论了无数个问题之后的笔记式总结。主要聊聊三个东西Pulsar的架构逻辑为什么在2025年更有说服力、现场被问爆的“页面里怎么看消息”到底怎么解决、“延迟消息队列”在Pulsar里是怎么实现的。无论你正在做MQ选型还是已经上了Pulsar但还没吃透原理这篇都值得花十分钟读完。1. Make MQ Great Again背后的现实Pulsar为什么值得被重新认真对待1.1 从一句口号看社区的心态变化“Make MQ Great Again”这句口号与其说是挑衅不如说是社区给自己打气。现场一个分享者说得很直接消息队列这十年不缺性能数据缺的是“用起来不难受”。Kafka用高吞吐证明了自己但高吞吐之后的运维压力、扩缩容成本、跨地域复制复杂度也是真实存在的。Pulsar这波要做的不是再做一个Kafka而是把消息队列从“管道”升级成“数据基础设施”——数据存得下来、流量扛得起来、业务隔离得开、历史数据查得到。一个细节让我印象很深现场问卷里有个问题是“你现在用Pulsar最头疼的事情是什么”答案最多的不是“不稳定”也不是“资料少”而是“怎么把消息捞出来看看”。这说明Pulsar的用户已经跨过了“能不能用”的阶段进入“好不好用”的阶段。1.2 开发者日的议程没有飘都在讲“用起来”这次Pulsar Developer Day放在COSCon25中国开源年会的大框架下分享密度不低。有讲存储引擎底层的有讲迁移实战的有讲流处理和函数计算的还有一整节是社区PMC在线答疑。最让我舒服的是大部分分享都不是论文式的“我们做了一个系统”而是“我们踩了这些坑、最后怎么绕过去的”。台下提问也在验证这一点——问得最多的是订阅模型怎么理解、Topic层级怎么规划、延迟消息为什么积压在Backlog里却不消费。这其实是生态成熟的表现。前几年大家问“Pulsar是什么”现在问“Pulsar怎么调”。从“What”到“How”中间隔着的正是大量真实的生产案例。2. 贯穿全场的架构主线存算分离到底解决了谁的痛点2.1 存算分离不是架构噱头是运维和成本问题Pulsar和Kafka最根本的区别就是Broker不存数据数据落在底层的BookKeeper集群里。这个设计听起来抽象但类比一下就很好懂Kafka像一个自带保鲜层的冰箱鱼肉虾蟹都塞在同一个机器里想多加一个抽屉得把整个冰箱搬回家Pulsar则是中央冷库加多个取货架货架只放近期要用的热数据大部分冻品在冷库里取货架不够了就多摆几个冷库不够了再单独扩建。现场有位从Kafka迁过来的运维同学分享说最大的感受是“终于不用盯着磁盘水位了”。在Kafka里分区扩容和数据均衡是联动操作节点磁盘快满了就得小心翼翼地迁移副本在Pulsar里Broker层是无状态的流量上涨时加机器就行存储层BookKeeper也可以独立扩缩容。这个解耦对运维来说不是“优雅”这么简单而是深夜被报警电话叫醒的概率直线下降。当然存算分离不是没有代价。数据在Broker和Bookie之间多走了一层网络小集群、低流量场景下反而会增加延迟。选型的时候要心里有数如果业务量级很小、峰值不高单体架构的MQ可能更省心存算分离的红利在集群规模和流量波动达到一定量级之后才会显著体现出来。2.2 多租户与IO隔离同一个集群里跑“不同脾气”的业务Pulsar的命名空间是三层结构Tenant租户→ Namespace命名空间→ Topic主题。这个设计最直接的好处是一个集群可以同时给多个业务线共用互相之间在权限、配额、备份策略上都能隔离开。会上有个案例很有代表性同一套Pulsar集群一个命名空间跑在线交易消息另一个命名空间跑离线数仓的批量导入。离线任务经常一次性灌上百万条消息如果是单体MQ很容易把同一批Broker的带宽打满在线业务跟着遭殃。Pulsar的做法是从Namespace维度做隔离策略配合独立Bookie集合把“爱折腾”的业务和“怕打扰”的业务分开在线交易消息的P99延迟几乎不受影响。我自己的经验是很多团队刚上手Pulsar时会忽略Namespace规划所有Topic堆在同一个默认Namespace里过了半年发现权限想分分不开、流量想隔离没得隔离。磨刀不费砍柴工第一天就把Tenant和Namespace按业务线规划好后面能省很多事。2.3 分层存储历史数据不再是一笔只能“删”的账Pulsar的分层存储Tiered Storage是我个人最看重的功能之一。消息在BookKeeper里保留一段时间后可以自动卸载到对象存储S3、GCS等本地只保留热数据。这意味着同一份数据可以“无限期”保留而不必为历史消息疯狂加盘。现场分享里有个数据他们保留了两年全量消息存储成本比本地SSD方案降了约八成查询历史消息时从对象存储拉取耗时会高一些但一天前到三个月前的热数据都在本地缓存真正需要翻几个月前消息的场景并不多。这个账算下来非常划算。不过要注意分层存储是“冷热”不是“快慢”。如果业务有高频查询历史消息的需求得评估好对象存储的读取延迟是否能接受别把冷数据当热数据用否则应用体验会打折扣。3. “页面里怎么看消息”这个高频问题背后的工具链真相3.1 先说结论Pulsar Manager能看什么不能看什么“mq怎么在页面查看消息”是这次大会上被追问最多的问题之一也是很多刚接触Pulsar的人最先找的东西。大家习惯用Kafka的UI工具看消息或者用RocketMQ的Dashboard到了Pulsar这边会自然地问我的Topic里现在有哪些消息内容是什么官方提供的Pulsar Manager确实是个Web管理界面但它更偏管理面查看Topic列表、订阅状态、Backlog积压、消息速率、连接数等。如果只是“看状态”Pulsar Manager完全够用但如果你想“看某条消息的内容”情况就麻烦一些——不同版本对消息内容的支持不一样而且消息查看会受到订阅模型的影响搞不好会干扰生产消费。这里有一个非常重要的原则生产环境上不要随便用一个普通订阅去“读消息”因为读取会推进游标可能把这条消息消费掉。要安全地看消息必须用独立订阅名或者用Reader模式。这个原则是所有“页面看消息”操作的地基。3.2 命令行三板斧admin看状态client读消息REST救急如果你记不住各种页面工具先把这三条命令吃透绝大多数“查消息”的需求都能解决。第一条看Topic的整体状态包括订阅列表、Backlog数量、每个订阅的各消费者状态pulsar-admin topics stats persistent://public/default/orders第二条用一次性订阅读取消息内容pulsar-client consume persistent://public/default/orders \ --subscription-name inspect_sub \ --num-messages 10这个操作会创建一个临时订阅“inspect_sub”读10条消息。因为是独立订阅名不会影响线上业务在用的订阅。读完之后如果想清理记得把临时订阅删掉pulsar-admin topics unsubscribe persistent://public/default/orders \ --subscription-name inspect_sub第三条通过REST接口直接看管理信息便于脚本化curl -s http://localhost:8080/admin/v2/persistent/public/default/orders/subscriptions对比下来页面工具适合“人肉监控时要个全局视野”命令行适合“我要精确知道某一条消息长什么样”。两者配合基本能覆盖大部分场景。3.3 页面化查看消息的轻量土办法如果你就是想要一个页面让团队同事不用敲命令就能看消息我推荐一个轻量方案基于Reader模式封装一个只读查询服务。Reader模式与订阅无关不推进任何游标不会把消息“消费掉”非常适合做消息体检和故障排查。大体思路是后端接口先通过getLastMessageId()拿到最新一条消息的MessageId再按需向前翻页把消息体反序列化后给前端渲染。前端用Grafana或者一个简单的Web页面就行。这样既能实现“页面看消息”又不用担心影响生产订阅。我自己给团队搭过类似的工具大概两百行Java代码就够用远比想象中简单。注意Reader模式虽然安全但读取大量历史消息会占用Broker和存储资源别拿它当全量导出工具控制好查询窗口。4. 延迟消息不是“定时睡觉”时间轮、索引与投递精度的工程取舍4.1 先理清概念延迟消息到底在延迟什么“mq延迟消息队列”这个热词背后是一类非常常见的业务需求订单支付超时后自动关单、优惠券到期前提醒、定时任务分批次执行、下单后30分钟未支付推送一条挽留消息。传统做法是建一张任务表起个定时任务轮询或者用Redis的ZSet按到期时间排序扫描。这些方案在数据量小的时候没问题但消息量一上来就成了“能用但不敢拆”的定时炸弹。消息队列原生的延迟消息核心思路不是“把消息晚点发出去”而是“消息已经写进队列了但不到时间不让消费者看见”。这个区别很关键消息是持久化的、有副本的、不会因为进程重启就丢延迟期间索引在Broker内存里到了时间再切回正常投递流程。4.2 Pulsar的时间桶设计从优先队列到更省内存的演进Pulsar在Broker端有个专门的组件叫DelayedDeliveryTracker负责管理延迟消息的触发。早期版本用的是一个内存优先队列按到期时间排序存消息位点实现简单但延迟消息量一大的时候内存占用和排序开销都很可观。社区在2.4之后引入了基于时间桶的默认实现把时间划分成一个个桶桶的宽度由Tick时间决定比如默认1000ms一个格。每条延迟消息根据到期时间落入对应桶桶内维护消息位点列表。时间轮走到的当前桶会把到期消息批量交给Dispatcher投递给消费者。这个“批处理”的优化思路和很多定时任务框架是一脉相承的不是每条消息一个定时器而是把相同时间窗内的消息攒在一起统一触发。4.3 一个延迟消息的完整旅程从Producer到Consumer以最常见的“订单超时30分钟关单”为例Java端的生产消息可以这样写producer.newMessage() .value(order-123456-timeout.getBytes(StandardCharsets.UTF_8)) .deliverAfter(30, TimeUnit.MINUTES) .send();代码很简单但背后Broker做了几件事Producer把消息发送到BrokerBroker写入BookKeeper保证持久化。写入成功后Broker的Dispatcher发现这条消息的DeliverAtTime还在未来就不把它放入正常可投递队列而是把消息位点交给DelayedDeliveryTracker。Tracker按到期时间组织索引时间轮在每次Tick时检查有没有到期批次。到期后Tracker把消息位点交还给DispatcherDispatcher再按正常流程投递给Consumer。整条链路里最容易让人误解的地方是延迟消息在到期前会以Backlog的形式显示在Topic统计里但消费者就是收不到。这不是Bug而是“消息已经入库、但投递被故意扣住了”的正常表现。4.4 精度、容量与订阅模式的取舍现场有好几个人问同一个问题Pulsar延迟消息的精度能不能做到毫秒级答案是可以但不建议。Broker里有个参数叫delayedDeliveryTickTimeMillis默认是1000ms也就是说到期时间会按1秒的粒度取整触发。把参数调小到100ms甚至10ms能提高触发精度但代价是时间轮桶数量变多、检查频率变高CPU和内存开销都会上升。绝大多数业务场景1秒精度完全够用真的要做毫秒级定时不如放到业务代码里去算。容量方面更要注意延迟消息的索引在Broker内存中如果系统里有上千万条延迟几个小时的订单消息内存压力会很大。社区版本在高内存压力时可能采取降级策略延迟消息被提前投递。所以设计时一定要评估“同时处于延迟状态的消息总量”别把长时间大范围的延迟任务一股脑塞进MQ。另外如果订阅模式用的是Key_Shared延迟消息的兼容性要提前验证。现场有个同学提到他踩过这个坑Key_Shared下同Key的延迟消息会影响后面紧挨着的非延迟消息的投递。遇到这种场景可以考虑换用Shared订阅或者在业务层把延迟消息和非延迟消息拆到不同Topic。方案精度持久化复杂度适用场景定时任务轮询数据库表分钟级依赖数据库低小规模任务Redis ZSet扫描秒级可能丢中延迟任务量中等Pulsar原生延迟消息秒级可调跟随消息持久化低大规模延迟消息自研时间轮毫秒级需自行设计高极高性能要求场景5. 现场排障案例与Demo教训比议程更值钱的东西5.1 消费者“不消费”一查Receiver Queue发现消息全在某一个客户端里现场一个工程师问为什么我有个消费者订阅了消息但Backlog就是不下掉排查过程很有代表性。先看Topic统计消息确实已经分发给了某个消费者但那个消费者迟迟没有Ack。再看客户端日志发现该消费者配置的ReceiverQueueSize是10万消息一次性全拉到了本地内存客户端处理线程跟不上堆积在本地Queue里Broker自然认为“消息已经投递成功”不会再重发给其他消费者。这个问题的本质是消息“从Broker视角已经发出去了但从业务视角还没处理完”。解决办法是控制单客户端拉取量把ReceiverQueueSize调到合理范围同时开启maxUnackedMessagesPerConsumer限制未Ack数量避免单消费者拖垮整个订阅组。5.2 订阅重建后消息从最早开始重复消费游标被重置了另一个参会者遇到的场景是他删掉了一个订阅重新创建了同名订阅发现消费者一启动从几周前最早的消息开始疯狂消费差点把下游服务打爆。原因其实不复杂——新订阅重新创建时游标位置取决于subscriptionInitialPosition参数。如果代码里设置成了Earliest新订阅就会从最早的消息开始读如果设置成Latest默认就是Latest则会从创建时刻之后的新消息开始读。这个坑的隐蔽之处在于同一个订阅名删掉再建看起来“还是那个订阅”实际上游标已经重置了。排查时一定要确认是“沿用旧游标”还是“重新创建订阅”不要想当然。生产环境中对订阅做删除和重建操作前最好先做好确认否则下游会被历史消息瞬间灌满。5.3 Demo翻车现场默认5MB消息上限很多人第一次知道下午有个演示环节展示消息“大文件”发送现场直接报错消息大小超过限制。台下一阵安静然后好几个人拿起手机拍照。Pulsar默认的单条消息大小上限是5MB超出会被Broker拒绝。这个限制在官方文档里写得清清楚楚但说实话不踩一次坑很难记住。更值得讨论的不是怎么调大上限而是很多人压根没想过“大对象应该放哪”。经验做法是超过几百KB的二进制内容先存到对象存储消息里只放引用地址。这样既绕开了消息上限也避免了Broker处理大Payload带来的内存压力。如果确实需要调大可以修改Broker配置文件里的maxMessageSize参数但一定要先评估内存和网络带宽。5.4 给正在选型或刚上手的人几句实在话活动散场前我问了自己一个问题如果要给刚接触Pulsar的人三个建议我会说什么。第一先用Docker搭一个Standalone实例亲手把Tenant、Namespace、Topic三层关系建一遍很多概念就通了。不要一上来就照着博客抄配置先亲手创建几个Topic感受一下“同一个集群里可以有多套隔离空间”是什么意思。第二别把Kafka的经验原样搬到Pulsar上。消费位移的管理方式不同订阅模型的灵活性不同连“同一个Topic可以同时有多个互不干扰的订阅”这种基本特性很多人都是用了半年才反应过来。第三生产环境一定要熟悉pulsar-admin topics stats-internal这类底层命令输出的字段含义尤其是PendingAcks和Cursor相关的信息——这些才是你半夜排查问题时的真正抓手。散场之后走出会场的时候天已经有点暗了。屏幕上“Make MQ Great Again”的投影早就撤掉但那个梗带来的讨论还在耳边。MQ伟大不伟大不是一句口号能决定的也不是某个社区在台上喊几声就行的。它取决于每天凌晨两点还在盯Backlog的运维取决于那些在生产环境里把消息队列调出低延迟、高吞吐的工程师也取决于每个愿意把“用起来不难受”当目标去改进的贡献者。这次Pulsar Developer Day给我最大的收获其实不是哪个具体的架构方案而是看到有一群人真的在认真解决“消息队列不止要快还要好用”这件事。延迟消息怎么更合理地扣住、消息怎么被安全地查看、多租户怎么隔离得干净——每一个问题背后都有人在踩坑、修补、总结再写成文档分享出来。技术也就在这一次次“看得见的麻烦”里往前走了一步。
返回列表