ARTICLE DETAIL

资讯详情

深耕郑州网站建设与运营推广的一线实战洞察。

Pulsar实战回顾:从Kafka迁移到MessageId深度解析

Pulsar实战回顾:从Kafka迁移到MessageId深度解析 刚从广州回来COSCon‘25和Pulsar Developer Day 2025同场举办两天听下来我最大的感觉是消息队列MQ这个存在了十几年的“老家伙”正在以一种很微妙的方式重新成为架构圈的主角。Pulsar的专场几乎坐满了人现场提问集中在同一个方向“我到底该不该从Kafka迁到Pulsar”以及“Pulsar的MessageId为什么长成这样”这篇文章不打算复述所有演讲PPT我想把现场听到的、亲手试过的、回来又踩过坑的内容揉在一起做成一份可以直接拿去用的Pulsar实战回顾。无论你是在做技术选型还是刚被一个莫名其妙的MessageId报错卡住这篇都值得你花十分钟看完。1. COSCon‘25现场的MQ氛围Pulsar凭什么站在聚光灯下1.1 一场会议把“消息队列”重新拉回话题中心COSCon中国开源年会本身是一年一度的开源大聚会今年正好和Pulsar Developer Day 2025合办所以会场上多了不少做消息中间件、流处理、事件驱动架构的开发者。我在展区逛了一圈发现很多项目已经不再用“我用了Kafka”来介绍自己而是改说“我们用Pulsar做了一个多租户事件平台”。这个变化很微妙说明业界对MQ的需求已经从简单的“削峰填谷”升级成了“数据基础设施”。而且这次会议把MQ放到了一个特别显眼的位置专门有半天是Pulsar的开发者日议题覆盖存储层、协议兼容、生态集成、生产环境运维。和几年前“MQ就是发消息、收消息”的讨论完全不一样了。给我的感觉是消息队列正在从“中间件”变成“平台底座”而Pulsar恰好是这个转变里最激进、也最完整的一个答案。现场还有一个很有意思的细节很多讲师在讲架构演进时都用“从单体到微服务再到事件驱动”作为引子几乎每个人都会提到“你需要一个能支撑多年业务增长的消息系统”。这不再是Kafka一统天下的叙事而是Pulsar和Kafka在架构理念上的正面碰撞。1.2 MQ在云原生时代到底变成了什么角色过去我们聊MQ第一反应是拿来解耦和异步订单系统发一条消息积分系统慢慢消费谁也不用等谁。但COSCon’25上大家聊的MQ早就超出了这个范畴。它现在是流处理的数据源是微服务之间的“神经系统”是多云环境下的统一数据通道甚至是数据库变更捕获CDC的输送管道。这种角色变化对技术栈提出了两个硬性要求一是要能弹性扩展二是要能长期低成本保存数据。Kafka在很长一段里靠分区和副本解决了这两个问题但遇到集群规模变大、Topic数量变多时运维成本会直线上升。Pulsar在这个背景下“翻红”原因就是它把存储层单独拎出去用BookKeeper做了一套可以独立扩缩容的架构。说白了Pulsar不是在跟Kafka抢“消息引擎”的位置而是在定义一个更符合云原生思路的“消息存储平台”。我在现场听一位从Kafka迁到Pulsar的工程师分享他说了一句话我印象很深“Kafka像是一个把所有功能都揉在单车上的方案你骑到一定速度就得换发动机Pulsar更像是把发动机和车轮拆开卖你可以单独升级发动机也可以只换轮胎。”这个类比虽然糙但一下子把存储与计算分离的好处说清楚了。1.3 现场最让我意外的几个Pulsar话题第一天主会结束后Pulsar专题区有几个话题是我完全没想到的。第一个是“Pulsar on Kubernetes到底踩了多少坑”讲师直接晒出他们的内存参数和堆外内存配置听到一半我就觉得这趟来得值。第二个是“用Pulsar做事件溯源架构”这在国内其实还比较小众但现场提问的人特别多。第三个是关于“Pulsar协议兼容”的讨论官方说已经兼容Kafka协议很多人直接问“那我能把Kafka客户端原封不动接过来吗”答案是可以但有几个细节必须注意。这让我意识到一个开源项目的社区热度往往不是看它发了多少版本而是看有多少人愿意在生产环境里“趟雷”并把经验分享出来。Pulsar的社区明显在进入这个阶段。出门的时候我打开手机上之前保存的一个链接——mq官网上的快速入门突然觉得这次回去得好好把Pulsar的底层机制捋一遍不然对不起现场吸收的信息量。2. Pulsar核心亮点拆解不只是“又一个Kafka”2.1 存储与计算分离的架构到底解决了什么问题要理解Pulsar绕不开“存储与计算分离”这六个字。传统MQ比如Kafkabroker既负责接收消息、维护元数据也要负责把数据写到本地磁盘并做副本同步。这意味着如果你需要更大的存储容量通常只能给broker加磁盘或者加机器然后把分区重分配。Topic数量和分区数量一旦上去集群的元数据压力和不均匀问题就会开始冒头。Pulsar则把数据实际存储放到了BookKeeper集群里broker层只做消息的接收、路由和分发不直接持久化数据。写入的消息会被切分成段segment均匀分布到多个BookKeeper节点上读取时broker再向BookKeeper拉数据。这个拆法带来的直接好处是你想扩容存储就单独加BookKeeper节点你想提高吞吐就单独加broker节点。它们不会互相拖累。这就是为什么Pulsar可以做到“千万级Topic”而Kafka如果Topic过多分区的副本和元数据会让集群很吃力。从使用者的角度看你连接的是broker但消息实际落在BookKeeper里。所以排查问题的时候不能只盯broker日志还要看BookKeeper磁盘和写入延迟。现场有位嘉宾分享过一个案例Pulsar吞吐突然下降最后发现是某台BookKeeper的磁盘写入延迟飙到了200毫秒而broker和客户端监控却表现正常。这就是存储与计算分离架构下独有的排查思路也提醒我们不要把Pulsar当成一个“黑盒”来用。2.2 Pulsar的MessageId为什么长这样以“messageid|28077:20854:0”为例很多刚接触Pulsar的人都会被MessageId吓到比如在控制台或日志里看到一串像“messageid|28077:20854:0”这样的字符串第一反应是“这是什么乱码”我第一次看到也懵了。其实这个格式一点也不复杂它就是Pulsar里消息位置的坐标拆开来看分段含义messageid固定前缀表示这是一个MessageId28077ledgerId也就是BookKeeper里这本账本的编号20854entryId消息在账本里的条目序号0partition index如果这条消息来自分区Topic这里就是分区编号非分区Topic通常为0为什么用ledgerId和entryId而不是像Kafka那样的单调递增offset因为Pulsar采用分段存储消息先落到一个ledger里ledger写满或达到阈值后会关闭再创建新的ledger。所以“28077:20854”描述的其实是消息在BookKeeper中的物理位置。这种设计好处是便于并行写入和快速恢复也让Pulsar天然支持从任意时间点、任意位置重新消费。现场有人问“这个ID是全局唯一的吗”答案是ledgerIdentryId在同一个Topic的同一个分区内可以唯一定位一条消息但如果你要跨Topic对比就必须连同Topic名一起看。还有个小知识点在客户端API里MessageId通常是个Java对象它的toString()输出可能是类似“28077:20854:-1”的格式而我们在日志或工具里看到的“messageid|28077:20854:0”很可能是某个封装层加了前缀。如果你需要手动解析按竖线和冒号分割就能拿到三个数字别被格式吓退。2.3 消息确认、延迟队列、分层存储等易被忽略的细节除了MessageIdPulsar里还有几个概念特别容易被忽略但又是生产环境离不开的。第一个是消息确认ack机制。Pulsar的ack是按MessageId精确确认的消费者可以逐条确认也可以累积确认。很多人以为确认之后消息就删了其实并非立即物理删除。每个订阅游标cursor会记录已经确认到的位置只有所有订阅都确认过了底层数据才可能被回收。这就意味着如果有一个消费者一直不消费消息会一直保留磁盘占用不会因为你其他消费者已处理完就释放。第二个是延迟队列。Pulsar从2.4版本开始支持延迟消息投递底层实现很有意思它并不是起一个定时扫描线程而是把延迟消息放到一个内部的重试主题通过时间戳来计算投递时间。你在生产端调用deliverAfter()或者deliverAt()看起来像发了一条普通消息实际Broker会把它暂存在一个特殊位置时间到了再转到真实订阅里。我现场听到有人说“Pulsar延迟队列不如RocketMQ好用”其实是因为没理解它的设计思路只要设好消费者ack超时时间体验并不会差。第三个是分层存储。Pulsar允许把旧数据从BookKeeper卸载到S3或其他对象存储broker在消费时再按需拉取。这个功能对超长保留时间的场景非常有用很多人并不知道。我在实操时最常遇到的问题是分层存储的配置文件写错导致数据无法读取所以提醒一句如果启用了tiered storage务必保证对象存储的bucket生命周期规则不会自动删除你的分块文件。3. 现场实操复盘从0到1跑通一个Pulsar生产级Demo3.1 环境准备与集群搭建在大会的Pulsar Developer Day上官方给我们准备了一个快速实验环境但我更建议你回到本地自己搭一遍因为只有自己从零搭过才会理解哪些是网络问题哪些是配置问题。最省事的方式是用Docker Compose跑一个standalone模式version: 3.8 services: pulsar: image: apachepulsar/pulsar:3.3.0 container_name: pulsar ports: - 8080:8080 - 6650:6650 command: bin/pulsar standalone启动后可以通过http://localhost:8080/admin/v2/clusters检查控制台是否正常空列表就说明服务起来了。生产环境当然不能用standalone但用它来学习MessageId、验证客户端行为完全够用。我在现场看到好几个人卡在“Pulsar明明启动了却连不上”八成是没等broker完成初始化默认standalone启动需要几十秒别急着马上发消息。如果你想更接近生产环境的形态可以起一个真正的集群至少需要3个BookKeeper节点、3个Broker和1个ZooKeeper或元数据服务。这里我不展开所有yaml只提醒一个最容易被忽略的参数managedLedgerDefaultMarkDeleteRateLimit它控制游标确认后删除消息的速率。如果这个值太低你会发现消费者已经确认了消息但磁盘占用还是不开心的往上涨。3.2 生产者、消费者与MessageId的观察最简单的Pulsar客户端用Java写起来非常直接PulsarClient client PulsarClient.builder() .serviceUrl(pulsar://localhost:6650) .build(); ProducerString producer client.newProducer(Schema.STRING) .topic(persistent://public/default/demo) .create(); producer.send(hello from COSCon25);发送成功后的返回值就是一个MessageId。如果你把它打印出来大概长这样MessageId msgId producer.send(hello); System.out.println(msgId.toString()); // 输出可能是28077:20854:0这个输出和我们在日志里看到的“messageid|28077:20854:0”是对应的只是少了前缀。如果你想拿到ledgerId和entryId的单独值可以用long ledgerId msgId.getLedgerId(); long entryId msgId.getEntryId(); int partitionIndex msgId.getPartitionIndex();我在现场演示时发现很多人有个误区认为同一个Producer连续发送的两条消息它们的ledgerId一定相同。其实不一定。当ledger写满比如达到managedLedgerMaxEntriesPerLedger设置的条数或者Broker要求滚动时Producer下一条消息就会进入新的ledger。所以MessageId里的ledgerId会跳动是正常现象不用紧张。消费者的核心代码也很简单ConsumerString consumer client.newConsumer(Schema.STRING) .topic(persistent://public/default/demo) .subscriptionName(my-subscription) .subscribe(); while (true) { MessageString msg consumer.receive(); System.out.println(msg.getMessageId().toString()); consumer.acknowledge(msg); }每次循环收到的消息其MessageId就是生产者写入时生成的同一个ID。这里要特别提醒如果消费者没有调用acknowledge并且启用了ackTimeout消息可能会被重新投递。消息重投时Pulsar会重新构造一条相同数据的消息但底层的数据物理位置可能不变MessageId依然相同。这会导致你的消费逻辑必须保持幂等不能光靠MessageId去重。3.3 压测参数与性能调优心得现场主办方安排了一个小压测环节Pulsar自带一个pulsar-perf工具用法非常简单bin/pulsar-perf produce -r 10000 -s 1024 persistent://public/default/perf-topic这个命令表示每秒发送10000条、每条消息1KB。我在现场跑下来standalone模式大概能支持每秒3-4万的写入但瓶颈很明显在磁盘IO。如果你是本地测试可以把managedLedgerDefaultEnsembleSize、managedLedgerDefaultWriteQuorum、managedLedgerDefaultAckQuorum都设为1省去副本同步的开销但生产千万别这么干。性能调优方面现场压测专家提到几个最容易踩的坑Producer端batchingMaxMessages默认是1000如果你发送的batch太小吞吐上不去建议结合消息大小调整到1000或2000。Consumer端receiverQueueSize默认是1000如果消费者消费逻辑较快但网络延迟高可以调大这个值。如果消息大小不一致启用compressionType如LZ4、ZSTD能明显降低带宽占用但CPU消耗会增加需要权衡。我把现场的压测数据整理了一下在4核8G的单机上standalone模式下接近“单生产者-单消费者”的场景大概能达到每秒4万条左右的稳定吞吐消息大小1KB。如果启用副本为3的集群吞吐会下降一半以上但这正是为了保证持久性需要付出的代价。压测过程中我还发现一个很有意思的问题消费者收到的MessageId并不总是严格的连续递增。因为Pulsar的消息可以分布在多个ledger里ledger间的切换会让entryId从0重新开始。所以你看到“28077:20854”之后下一条可能是“28078:0”别误以为丢消息了。整体连续性看ledgerId和entryId的组合不能只看entryId。4. Pulsar和Kafka选型资料丰富度之外我们更该看什么4.1 “Pulsar和Kafka哪个资料更丰富”的真实答案网上搜索“pulsar和kafka哪个资料丰富一些”你会得到一堆帖子。作为一个两边都写过生产代码的人我的结论很直接如果只看资料数量Kafka赢毕竟它火了快十年教程和踩坑帖子遍地都是。Pulsar的资料虽然在快速增长但深度实战文章还是少得多尤其中文社区里“读完能直接照抄”的内容更稀缺。但这不意味着Kafka资料丰富就是选它的理由。我发现一个规律资料越多信息噪音也越大。Kafka相关内容里大量都是“Kafka入门”“Kafka安装”之类的基础贴真正讲清“副本同步机制”“消费者Rebalance内幕”的反而要花时间去淘。Pulsar因为比较新能搜到的内容大多来自官方文档和少数有质量的技术博客信息密度反而可能更高。我在学习Pulsar时官方文档和几位核心贡献者的文章给到的帮助远超在搜索引擎里翻几十页Kafka旧帖。所以我的建议是理解技术选型时资料丰富度不是一个核心指标。你需要的不是“哪个能搜到更多”而是“哪个能让你在踩坑的时候更快找到答案”。Pulsar的社区响应速度和官方Slack确实很活跃这个价值是百度/谷歌搜索数量给不了的。4.2 选型考量生态、团队、运维成本在COSCon’25的圆桌论坛上主持人问了一个很尖锐的问题“你们为什么从Kafka迁移到Pulsar”现场举手的几个嘉宾回答并没有人说是“因为Pulsar资料更丰富”而是集中在三点需要支持大量Topic几千到几万Kafka分区多了以后Broker内存和元数据压力明显。需要多租户隔离不同业务部门用同一个集群但彼此资源要隔开。希望订阅模型更灵活能兼容流式处理和队列消费两种模式。这三点恰好是Pulsar的强项。Kafka的优势在于生态成熟、流处理Kafka Streams/KSQL一体性更好Pulsar的优势在于架构更有弹性Subscription模型支持exclusive、shared、failover、key_shared四种模式团队内部想按队列还是按流使用都行。运维成本方面Pulsar的存储与计算分离意味着生产环境至少要维护两套集群Broker和BookKeeper第一眼看上去比Kafka复杂。但Pulsar提供了非常清晰的管理命令和告警指标尤其pulsar-admin工具上手之后会觉得很多操作比Kafka的kafka-topics.sh更好用。选型时最忌讳“看着哪个火就选哪个”一定要先答清楚一个问题你的核心痛点是吞吐、Topic数量、多租户还是流处理API痛点决定选型。4.3 给初学者的学习路径建议如果你决定学Pulsar我建议别急着去查“Pulsar和Kafka哪个资料丰富”而是按下面这条路径走能少走很多弯路第一步跑一个standalone实例用命令行工具发送和消费几条消息感受一下Topic、订阅、MessageId的形态。第二步去官网mq官网以及Apache Pulsar官网看“Concepts and Architecture”一章把Ledger、Subscription Cursor、Managed Ledger这几个基础概念弄懂。这一步很重要很多人在没搞懂架构的情况下直接看客户端API看两天就放弃了。第三步自己写一个Java或Python客户端Demo分别验证“shared订阅下消息如何分发”“key_shared如何保证同一key有序”“延迟消息和重试队列怎么用”。建议把接收到的MessageId全部打印出来你会对Pulsar底层存储产生真实的体感。第四步模拟生产环境部署一个3节点BookKeeper 2节点Broker的集群。不用上太复杂的自动化用Docker Compose即可。这个过程遇到的问题会让你对Pulsar的敬畏感上升一个级别。第五步再去看社区里的生产实践分享这时候你看文章的速度和理解深度都会完全不一样。到了这一步你会发现资料丰富不丰富真没那么重要。5. 参会避坑指南与常见问题速查5.1 现场演示翻车问题大会上最不缺的是“演示Demo翻车”的瞬间。这次Pulsar专场也发生过几次最常见的是这三类standalone服务没等就绪就执行命令报Lookup failed错误。解决方法是等待日志出现“Messaging service is ready”再操作或者用bin/pulsar-admin namespaces list public验证。用Kafka客户端连Pulsar时忘了配置topic映射或协议端口。Pulsar兼容Kafka协议是通过Kafka protocol handler实现的需要把kafkaProtocol相关配置打开而且端口默认是9092而不是6650。消费者订阅了分区Topic但消息集中在某个分区导致有些消费者空闲。这个其实是因为生产端没有给key或者分区策略不合适Pulsar的RoundRobinPartition默认行为在某些版本下会让消息分布不均匀需要自己验证。这些坑并不是Pulsar独有的任何技术都有类似问题。但现场看到那么多人都栽在同一个地方说明官方文档里的“快速开始”其实跳过了很多前置细节。5.2 生产环境常见问题速查表我把这次会上听到的、加上之前自己踩过的Pulsar生产环境问题整理成一张表方便收藏问题现象可能原因排查建议消息积压但消费者正常订阅模式为shared且消费者数量少或ackTimeout过短导致重复投递检查订阅模型确认receiverQueueSize和ackTimeout设置Topic数量多后broker内存飙升Topic分隔成大量Fragment缓存占用过高调小managedLedgerCacheSizeMB限制缓存磁盘占用持续上涨有订阅游标未推进消息无法删除查询订阅积压移除失效的消费者或调大markDeleteRateLimit客户端报“Producer blocked”达到maxPendingMessages上限开启批量发送或调大该参数生产建议至少2000MessageId出现“”分隔且带前缀非标准中间件封装或日志框架序列化导致其中“消息积压但消费者正常”这个问题最容易在从Kafka迁移到Pulsar的团队里出现。Kafka的消费者是pull模型Pulsar虽然有类似机制但它的receiverQueueSize是broker主动推多少给客户端如果你消费者处理很快但网络带宽不够实际吞吐反而会被客户端缓存限制。记得多测压力值别只盯着CPU。5.3 我的几点经验心得参会回来后我又花了两个晚上把Demo重新跑了一遍有一些很个人、但大概率对你有用的体会第一Pulsar的学习曲线确实比Kafka陡一点。不是因为它难而是因为它多了一层BookKeeper的概念。如果你能接受“Broker负责干活BookKeeper负责存数据”这个世界观后面所有问题都会变简单。第二官方文档里关于MessageId的介绍很简短但你把底层存储模型搞清楚之后就会发现这个ID其实非常优雅。它能精确表示一条消息在哪个ledger、哪个entry、哪个分区这让Pulsar在处理回溯消费、重置游标时非常便利。Kafka的offset是个相对序号还需要依赖broker把offset转换为具体文件位置两者的设计哲学完全不同。第三也是最重要的别被“资料少”拦住。Pulsar再新也是一套开源中间件基础网络与分布式的知识都是通用的。你如果熟悉Kafka迁移到Pulsar时真正要学的只是架构上的几个新概念而不是重新学一遍消息队列。我在现场看到很多Pulsar的入门者其实都是一边看官方文档一边就上手了。最后再分享一个小技巧如果你在消费端看到MessageId突然跳变不要慌张用我上面提到的getLedgerId()和getEntryId()分别打出来再结合时间戳就能判断是不是发生了ledger滚动。很多时候生产事故只是虚惊一场贴个数据到社区里问一句十分钟就能得到答案。Pulsar这个社区值得你去认真参与一次线下活动。
返回列表