每年开源大会扎堆的时候,各种同场技术专场就成了圈内人真正盯着的目标。今年 COSCon 2025 的 Pulsar Developer Day 议程一放出来,好几个做消息中间件选型的朋友就转给我看,问得最多的无非是几个老问题:Pulsar 到底比 Kafka 强在哪?网上资料怎么感觉没有 Kafka 多?还有人不理解收到的 Message ID 为什么是28077:20854:0这种奇怪格式。这篇文章就从开发者日和议程背后的技术趋势说起,把 Pulsar 的消息模型、Message ID 结构、快速上手路径和排查经验一次讲透,适合正在选型、刚接触 Pulsar 或者已经在生产环境踩坑的开发者参考。
1. 从开发者日看消息中间件的技术风向
1.1 为什么一个"同场活动"值得专门关注
很多人可能还没意识到,消息中间件这类基础组件,早就不再是"能用就行"的阶段了。像 Pulsar Developer Day 这种能进入 COSCon 主会场的开发者专场,本身就释放了一个信号:消息系统已经从后端工程师的私有话题,变成了整个技术圈都在关注的核心基础设施。你去现场听一圈就会发现,讨论的话题已经从"怎么搭个队列"变成了"怎么在多集群、多云环境下做到低延迟和高可用",这完全是两个维度的东西。
对开发者来说,这类同场活动的价值在于密度高。一天之内你能听到真实生产环境的踩坑案例、核心维护者讲设计取舍、不同的技术团队分享各自的风控和拆分方案。这比你自己翻几个月文档都管用,因为很多经验是文档里不会写的。举个例子,Pulsar 的存储层基于 Apache BookKeeper,这个设计解决了很多问题,但也带来了新的运维习惯要求,比如你不能再像看 Kafka 那样只看 broker 的磁盘水位,还要关注 Bookie 的存储均衡。这种信息,只有到了现场或者深入社区才能快速 get 到。
1.2 选型绕不开的话题:Pulsar 还是 Kafka
每次一聊消息中间件,Pulsar 和 Kafka 的对比就是必答题,这次的议程里也有不少内容是在回应这种选型焦虑。先说结论:两个项目都很成熟,没有绝对的"谁取代谁",关键是场景对不对得上。
Kafka 的优势非常明显:生态最完整、中文资料最多、很多团队从 0.8 时代就在用,踩坑经验前人早写完了。如果你只是做标准的日志管道、事件流分析,Kafka 几乎是无脑选择。它的问题在于:存储和计算是耦合的,分区数上去以后,broker 的迁移和扩容都比较重,多租户和跨地域复制也不是它的强项。
Pulsar 走的是另外一条路。它把 Broker 和存储层拆开,Broker 本身无状态,消息实际存在 BookKeeper 里,这个架构带来了几个实打实的好处:
- 扩容时不需要搬迁消息数据,新加 broker 就能接入流量
- 存储和计算独立扩容,存储量大的 topic 不会拖垮计算能力
- 多租户隔离做得更细,不同团队的 topic 可以设置不同的策略,互不干扰
- 跨地域复制是内置能力,对全球化业务非常友好
下面这个表可以比较直观地看出两者侧重点:
| 对比项 | Kafka | Pulsar |
|---|---|---|
| 存储与计算 | 耦合,分区迁移较重 | 分离,Broker 无状态 |
| 多租户 | 靠配额和 ACL 实现,相对基础 | 内置租户/命名空间,策略更细 |
| 跨地域复制 | MirrorMaker 等工具方案 | 内置复制机制 |
| 消息保留 | 基于时间/大小删除 | 支持分层存储和更灵活的保留策略 |
| 消费模型 | 基于 offset | 基于 Message ID 和 Cursor |
| 资料丰富度 | 非常丰富 | 官方文档系统,社区资料在快速增长 |
关于资料丰富度的问题,我的看法是:如果只看中文二手资料,Kafka 确实多,但 Pulsar 的官方概念文档写得非常清楚,架构白皮书也值得反复读。你只要把 Ledger、Cursor、Subscription 这几个核心概念弄明白,很多知识是可以从 Kafka 迁移过来的,并不冲突。
1.3 创新实践都在做哪几件事
从议程的主题分布来看,消息中间件目前的创新主要集中在三个方向。
第一是云原生和 Serverless 化。Pulsar 的存储计算分离架构天然适合往云上搬,很多团队在讨论怎样让 topic 按需扩缩容、怎样把资源成本做到按量付费。第二是多集群容灾和一致性。业务跨地域部署之后,消息系统怎么保证数据不丢、顺序可控,同时还要在机房故障时快速切换。第三是成本和可观测性。消息系统跑久了,磁盘和带宽成本会非常扎眼,怎么通过分层存储、压缩策略来控制成本,同时又能让消息的流转状态被准确追踪,这些都是真实痛点。
议程里有人分享 Topic 数量从几百涨到几万之后的治理经验,也有人讲如何把消费延迟控制在毫秒级,这些内容背后其实都在解决同一个问题:当消息系统成为业务主动脉之后,稳定性和成本之间怎么找到最优解。
2. Pulsar 核心原理与 Message ID 深度解析
2.1 三层架构:Broker、BookKeeper、ZooKeeper 各司其职
很多刚接触 Pulsar 的人会被它的组件数量吓到,觉得比 Kafka 复杂多了。其实拆开看逻辑非常清晰。Pulsar 分成了三层:最上面是无状态的 Broker,负责接收生产者和消费者的连接、处理各种协议;中间是持久化存储层,由 Apache BookKeeper 提供,负责真正把消息落盘;最下面是元数据服务,传统上用 ZooKeeper 来管理整个集群的状态,比如 topic 分布在哪些 broker 上、bookie 是否在线。
为什么要把存储单独拆出来?这个设计我打个比方你就明白了。Kafka 相当于每个饭店自己建一个仓库,生意好了就得多建几个仓库,还要把食材搬来搬去。Pulsar 的 BookKeeper 则像一个中央冷库,每个饭店(Broker)只负责接单和出菜,食材统一存在冷库里。这样饭店扩不扩张,跟冷库容量的关系就没那么大了,冷库里的食材还能被多个饭店共享。
具体到 BookKeeper,它的核心概念是 Ledger。一个 topic 的数据会被切分成多个 Ledger,每个 Ledger 里的数据条目叫 Entry。这个设计带来了一个直接好处:Pulsar 不用像 Kafka 那样需要顺序读写一整段的日志文件,它可以并行写多个 Ledger,扩容时也不用搬数据。你在 Pulsar 后台看 Topic 的读写速率时,偶尔会注意到 Ledger 切换,比如写满一个 Ledger 会自动创建下一个,这在底层是非常高频且自然的操作。
2.2 Message ID 为什么长这样:ledgerId:entryId:partition
现在来聊那个很多人问过的问题:为什么 Pulsar 的 Message ID 是28077:20854:0这种格式。先说结论:这个 ID 由三部分组成,依次是ledgerId:entryId:partition-index。拿你看到的例子来说,28077是这条消息所在的 Ledger 编号,20854是这条消息在 Ledger 里的 Entry 序号,最后的0代表它是这个主题的第 0 个分区。
理解这个 ID 的关键,是要忘掉 Kafka 里简单的"offset+partition"思维。Kafka 的 offset 是一个分区内从 0 开始递增的序号,你只要知道 offset 就能定位消息在日志文件中的大致位置。但 Pulsar 因为是存储计算分离,消息存在 BookKeeper 的 Ledger 里,光一个数字无法定位。它需要两个坐标:一个指定在哪个 Ledger,一个指定在 Ledger 里的哪一条 Entry。这就像你去图书馆找书,不能只说"给我第 100 层书架上的第 50 本书",你还得告诉管理员是哪个书库(Ledger),这样两个坐标唯一确定一条消息。
那个partition-index则对应分区编号。如果主题是非分区的,这个值通常是-1。如果是分区主题,第几个分区就会在这里体现。
再说一个容易踩坑的点:直接用 Message ID 做跨集群的判断没有意义,因为不同集群的 Ledger 编号体系是独立的。同一个 topic 在两个集群里的同一逻辑消息,Message ID 完全不同。所以在做容灾切换或者多集群同步时,不要依赖 Message ID 做全局唯一标识,该加业务主键的还是要加。
2.3 从 Message ID 到消费位点:Cursr 和订阅模型
Message ID 除了用来定位消息,还有一个重要作用:管理消费位点。Pulsar 里叫 Cursor,它记录的是当前订阅已经消费到哪个位置。和 Kafka 的消费组 offset 不同,Pulsar 支持一个主题上同时挂多个订阅,每个订阅可以有自己独立的消费进度。这意味着同一份数据可以按不同业务需求分别消费多遍,这在 Kafka 里需要做配置在主题上建多个消费组,但 Pulsar 把这件事做成了订阅层面的灵活机制。
订阅模型有四种,选错会直接影响使用效果:
Exclusive:独占订阅,一个 topic 同时只能有一个消费者在消费,适合严格顺序场景Failover:故障转移,多个消费者同时连接,但只有主消费者在处理,主节点挂了才切换Shared:共享订阅,消息在消费者之间轮询分发,处理能力强,但完全无序Key_Shared:按 key 分发,同一个 key 的消息永远进同一个消费者,兼顾顺序和并发
实际项目中,Shared 是很多团队默认的选择,因为它能最大化吞吐。但如果你对顺序要求高,比如一个订单的创建和取消消息必须被同一个消费者按序处理,那就要用 Key_Shared 或者 Failover 的独占语义。我见过不少线上问题,都是因为用了 Shared 订阅处理订单状态变更,导致同一条业务链的先后顺序乱了,排查起来非常痛苦。所以在创建订阅之前,就要想清楚顺序和并发之间的边界。
3. 快速上手 Pulsar 的实操路线
3.1 本地环境搭建与基础验证
如果只想在本地把 Pulsar 跑起来,最快的方式是 Docker。先拉镜像,然后以 standalone 模式启动,这个模式把 Broker、BookKeeper、ZooKeeper 都集成在一个进程里,方便开发调试。
docker pull apachepulsar/pulsar:3.3.0 docker run -it \ -p 6650:6650 \ -p 8080:8080 \ apachepulsar/pulsar:3.3.0 \ bin/pulsar standalone启动成功后,6650 是客户端连接的端口,8080 是 admin 接口。这时候可以开一个新的终端进容器,用自带的命令行工具做一次生产消费的验证。
# 进入容器 docker exec -it <容器ID> /bin/bash # 创建主题 bin/pulsar-admin topics create persistent://public/default/quickstart-topic # 生产消息 bin/pulsar-client produce quickstart-topic --messages "hello pulsar, from coscon" # 消费消息 bin/pulsar-client consume quickstart-topic -s "first-subscription"看到消息能正常收发,说明环境没问题。接下来可以拉一个 Python 客户端做简单的代码验证,Pulsar 支持的语言客户端很多,Python 是最容易上手的。
import pulsar client = pulsar.Client("pulsar://localhost:6650") # 生产者 producer = client.create_producer("quickstart-topic") producer.send(("hello from python writer").encode("utf-8")) # 消费者 consumer = client.subscribe("quickstart-topic", "first-subscription") msg = consumer.receive() print(msg.data()) print(msg.message_id()) # 看这里就能拿到类似 ledgerId:entryId:partition 的 ID consumer.acknowledge(msg) client.close()我建议你重点看一下msg.message_id()的输出,这能帮你直观理解上一节说的 Message ID 结构。第一次跑通了之后,再去读那些概念性的文档,会顺畅非常多。
3.2 生产与消费配置的关键参数
本地跑通只是第一步,真要上生产,有几个参数配置必须心里有数。
生产者的核心参数是batchingEnabled、batchingMaxMessages和pendingQueueSize。默认情况下 Pulsar 会开启批处理,把多条消息打包发送,吞吐会高很多,但延迟会略微增加。如果服务里有对延迟极敏感的业务,比如实时风控或交易回调,建议对相关主题关闭批处理,或者把批量大小调小。pendingQueueSize表示生产端在未收到确认前可以积压多少条消息,调太大会增加内存压力,调太小在高吞吐场景容易触发背压。
消费者的关键参数是receiverQueueSize和maxTotalReceiverQueueSizeAcrossPartitions。这个值表示消费者本地最多预取多少条消息。预取多了吞吐高,但消息堆积在客户端,服务重启时可能会有大量消息被重新拉取,造成重复消费;预取少了吞吐上不去,在高延迟网络环境尤其明显。通常在 1000 到 10000 之间根据实际压测结果来调。
还有一个容易被忽略的:ackTimeout和negativeAckRedeliveryDelay。生产环境偶尔会碰到消费线程卡死的情况,如果一直不 ack,Pulsar 会根据超时机制重新投递消息。超时设得太短,处理稍慢的消息会不断被重新投递,产生大量重复消费;设得太长,消费卡死时要等很久才会触发重试。我习惯把 ackTimeout 设为业务正常处理耗时的 5 到 10 倍,同时配合死信主题把多次重试仍失败的消息隔离出去。
3.3 主题策略与保留机制设置
Pulsar 的一个优势是可以针对命名空间或主题做细粒度策略管理。常见的配置有message-ttl(消息存活时间)、retention(消费者确认后消息保留时长)、backlog-quota(积压上限)。这个组合非常灵活,但也很容易误解。
举个例子,你设置message-ttl为 10 分钟,意思是这个主题里没有被消费的消息,10 分钟后会自动变成"已跳过"状态并对新消费者不可见。但如果你设置了retention策略,已经被消费和确认过的消息可以继续保留一段时间,供后续重新拉取或做数据回溯。这两个变量很多人会搞混,实际效果天差地别。
# 设置消息 TTL 为 1 小时 bin/pulsar-admin namespaces set-message-ttl public/default --ttl 3600 # 设置消息确认后保留 24 小时,并限制最大大小为 1GB bin/pulsar-admin namespaces set-retention public/default \ --size 1G --time 24h如果你有数据回溯的需求,建议把 retention 时间设置得比业务审计周期稍长一些。比如业务每个小时会跑一次数据核对,那就至少保留 24 小时,给排查问题留出余地。但如果是高吞吐的日志类主题,无脑保留大量数据会迅速吃光 Bookie 磁盘,这时候把 retention 设置为 0,只依赖 TTL 控制生命周期,反而是更合理的做法。
4. 高频问题排查与避坑指南
4.1 消息积压应该怎么查
消息积压是消息中间件最层出不穷的问题。我的排查顺序是固定的:先看积压量,再定位是生产快还是消费慢,最后找瓶颈。
在 Pulsar Admin 里可以直接看到主题的积压状态:
bin/pulsar-admin topics stats persistent://public/default/business-topic重点关注backlogSize和backlog两个字段。如果积压持续增长,先看消费者的并发数和处理耗时。如果消费者本身有外部 IO 或调用下游服务,通常是下游响应变慢导致消费整体拖慢。这时候最直接的办法是增加消费者实例,或者把部分流量切到新的主题。但注意不要只是盲目增加分区,如果下游数据库扛不住,加分区反而会把压力放大。
还有一种情况是生产端突然写入量暴增,导致消费端怎么都追不上。可以先看 topic 的msgRateIn和msgRateOut的差值。如果生产速率远高于消费速率,而且短时间内无法完成扩容,可以考虑先降级部分非核心业务对消息的依赖,比如把日志类消息暂时切到另一个低优路径,保证核心链路不丢消息。
4.2 重复消费与消息乱序
这是两个看起来类似但原因完全不同的经典问题。
重复消费的根源一般是 ack 丢失。消费者处理完业务逻辑,但还没来得及 ack,网络抖动或者进程重启,Pulsar 就会把这条消息重新投递给另一个消费者实例。这不是 Pulsar 独有的问题,任何消息系统只要禁止 at-least-once 语义都会引入这个现象。解法只有一个:消费端一定要做幂等。基于业务唯一键去重,而不是基于 Message ID 去重。因为同一逻辑消息重投时,Message ID 通常不变,你可以用 Message ID 做第一层拦截,但最终一致性还是得靠业务表里的唯一索引保证。
消息乱序的情况更多和订阅模型相关。如果你用Shared订阅,消息本身就是轮询分配的,顺序自然没保证。如果要用Key_Shared订阅来保证同一个 key 的顺序,要注意 key 的分布是否均匀。比如某个热点用户的订单量极大,hash 到同一个消费者后,其他消费者都在围观,总体吞吐反而会被这个热点约束住。这种时候只能做业务层面的拆分,让热点 key 进一步细分。
4.3 客户端连接与鉴权踩坑
很多团队第一次部署 Pulsar 集群时,会在鉴权配置上卡住。Pulsar 的鉴权体系和 Kafka 有些类似,但配置项更多。简单来说,你需要为 broker 开启认证,并配置授权设置。
# broker.conf 中相关配置 authenticationEnabled=true authorizationEnabled=true authenticationProviders=org.apache.pulsar.broker.authentication.AuthenticationProviderToken tokenSecretKey=...客户端连接时需要传入 token,否则会报AuthenticationError。有些朋友本地调试时明明没配鉴权,也会遇到连接被拒的问题,那通常是因为客户端用了pulsar+ssl://协议去访问普通端口。我用过一个比较省心的方式:先在本地全部用明文连接跑通功能,等要上生产前再统一处理 TLS 和 token,把协议换掉。这样能先把业务逻辑的问题排查完毕,再面对加密和鉴权的额外复杂度。
4.4 学习资料怎么看最有效
回到那个经常被问到的问题:Pulsar 和 Kafka 哪个资料更丰富,应该先学哪个。我的建议是,不要把它们对立起来。Kafka 的书籍和视频量大但质量参差不齐,很多讲的还是早期版本。Pulsar 的官方文档、架构设计文档以及 StreamNative 那套教程,相对更新,质量也更稳定。如果你完全没接触过消息中间件,可以先找一本 Kafka 的入门书建立基本概念,了解 topic、partition、consumer group,然后再看 Pulsar 里对应的概念是怎么设计的。这种对照学习法效率最高,因为很多名词你已经在 Kafka 那边理解了,换到 Pulsar 只需要关注它们"为什么不同"。
至于英文资料,Pulsar 官方博客、社区提案(PIP)和核心维护者的分享都值得读,尤其是 PIP,那是理解 Pulsar 后续演进方向的第一手来源。日常有问题先查官方文档,解决不了再去 GitHub Discussion,基本不会走偏。
回到 Message ID 那个疑问,当你理解了 Ledger、Entry 和存储分离架构之后,再看到那一串数字就不会慌了。它不是一个随机的哈希,而是消息在 BookKeeper 存储空间里的坐标。这个坐标让 Pulsar 可以做到独立扩容和精确回溯,是它相对传统队列的一大核心差异。对我个人来说,选型时真正让我倒向 Pulsar 的,并不是某个功能介绍,而是这些底层机制能不能让我在业务快速增长时不那么焦虑。消息系统的选型,往往到最后就是在选一种长期运维的确定性。希望这篇文章能把 Pulsar 最核心的概念和实操路径讲清楚,让你上手的时候少走几步弯路。