
看到 COSCon‘25 同场活动 Pulsar Developer Day 的议程正式发布我第一时间把整份议程从头到尾捋了一遍。Pulsar 这个项目从早期版本我就开始接触一路看着它从一个小众的消息队列变成云原生场景里绕不开的选项。这次开发者日把主题定为“聚焦消息中间件创新实践”既有架构层面的拆解也有落地过程中大家踩过的坑对于正在选型或者已经在生产环境里维护 Pulsar 的团队来说这份议程的信息量很大。消息中间件这几年越来越像业务系统的主动脉。以前我们聊消息队列无非是异步解耦、削峰填谷但现在的业务场景早就超出了这个范畴。事件驱动架构、实时数仓、IoT 数据接入、跨地域数据同步样样都离不开消息管道。Pulsar 能在这些场景里站住脚靠的不是简单的“又一个 Kafka”而是从设计层面把存储和计算拆开把多租户和跨地域复制做成内建能力。这篇文章我会结合议程的几个重点方向聊聊消息中间件在实践中的真正关注点同时把我自己从部署到排查的一些经验和教训一并写出来。1. 为什么消息中间件值得单独办一场开发者日1.1 从系统拆分到数据洪峰消息队列承担的角色变了早些年我们做系统拆分引入消息队列的核心动机很简单A 系统不需要同步知道 B 系统的处理结果中间塞一个队列生产者和消费者各自按自己的节奏走。后来双十一、秒杀这类场景把削峰填谷变成了刚需消息队列开始承担保护下游系统的职责。再往后业务开始要求实时性消息不仅要“能发能收”还要支持延迟消息、事务消息、死信重试甚至直接在里面跑轻量计算。我自己的体会是消息中间件正在从“技术组件”变成“数据基础设施”。你在上面流的每一条消息都可能触发下游的订单状态变更、风控模型更新、推荐特征拼接甚至驱动一个完整的业务闭环。这时候消息中间件就不再是“能用就行”而是要能保证顺序、能容忍突增流量、能跨机房同步、能隔离不同团队的业务数据。Pulsar 这轮议程里大量提到多租户、分层存储、跨地域复制正是顺着这条需求线展开的。1.2 Pulsar 与 Kafka 的差异不是“谁更好”而是“谁更合适”每次聊 Pulsar都绕不开 Kafka。很多团队选型的时候第一个问题就是“Pulsar 和 Kafka 到底差在哪”。我通常不建议别人用简单的一句话去概括因为这俩的设计哲学确实不一样。Kafka 是典型的分区和消费组强绑定模型写入和读取都围绕 partition 展开每个 partition 在存储上也是一份连续日志。这种设计让它在大规模日志管道场景下非常高效但对应的运维体验是 broker 有状态分区迁移和再平衡需要小心处理。Pulsar 则把 topic 的读写逻辑和底层存储拆开broker 本身不保存消息数据实际存储由 BookKeeper 负责。这带来的直接变化是 broker 可以更轻量地扩缩容存储和计算各自独立演进。这两种设计没有绝对的优劣更多是场景取向。如果团队已经在 Kafka 生态里沉淀了大量工具和运维经验继续用 Kafka 完全合理但如果业务对多租户隔离、跨地域复制、存储成本这些有更明确的需求Pulsar 的架构优势会逐步显现。我会在后面实测部分再展开聊先把概念理清看议程的时候才不容易被各种术语绕晕。1.3 开发者日为什么聚焦“创新实践”看议程标题就能发现这次 Pulsar Developer Day 并没有停留在“Pulsar 是什么”的基础科普上而是把重心放在了“怎么用好”和“怎么在复杂场景里落地”。原因我也能猜到Pulsar 本身的学习门槛比 Kafka 要高不少。Kafka 上手其实很快一个 producer、一个 consumer、一个 topic命令敲完就能跑。Pulsar 则多出 namespace、tenant、subscription、BookKeeper 这些概念第一次接触的人容易一头雾水。正因为门槛在这里社区里真正缺的不是文档而是“别人踩过的坑”。比如多租户的配额到底怎么设跨地域复制遇到网络抖动怎么处理存储积压导致磁盘写满该怎么办这些在官方文档里都有零散描述但没有经历过生产环境的人很难形成直觉。开发者日的 session 大多来自真实案例复盘会让这些隐蔽问题变得具体。我的建议是不管是刚入门还是已经深入使用都值得重点关注讲“故障处理”和“性能调优”的部分。2. Pulsar 核心概念梳理读懂议程必须补齐的背景2.1 分层架构Broker 和 BookKeeper 各管什么Pulsar 的架构可以用一句话概括计算和存储分离。客户端连接的是 brokerbroker 负责处理生产请求、路由消息、维护订阅游标等逻辑。但消息真正落盘的地方是 BookKeeper 集群它是一个低延迟的持久化存储系统内部以 segment 为单位管理数据。这个分层给我最直观的感觉是broker 变得“轻”了。以前用 Kafka 做扩容经常要担心 partition 迁移的过程中流量抖动、磁盘 IO 压力上升。Pulsar 里 broker 没有数据副本的压力加一个节点就能承接更多客户端连接和请求。BookKeeper 这边则是独立的高可用体系节点角色拆分更细读写路径经过精心设计哪怕是机械硬盘也能提供稳定的顺序写性能。不过分层设计也带来了额外的运维组件起码你要同时关注 broker 和 bookie 两套集群的资源使用情况。很多人第一次部署 Pulsar 时只盯着 broker 的 CPU 和内存忽略了 bookie 的磁盘空间和写吞吐这其实是一个很容易踩的坑。后面我会在排查部分再展开。2.2 订阅模型与消费语义Pulsar 的订阅模型和 Kafka 的消费组模型是两个需要刻意区分的概念。Kafka 里一个消费组在一个 partition 上只能有一个消费者消费者数量和 partition 数量需要精确匹配才能最大化吞吐。Pulsar 则把 topic 的消息流和消费进度拆成了订阅subscription这一层一个订阅下面可以有多个 consumer不同的订阅之间完全独立。Pulsar 提供了四种订阅类型独占、共享、故障转移、键共享。独占模式适用于必须严格有序且只能被单个消费者处理的场景共享模式适合吞吐要求高但不要强顺序的任务故障转移在正常情况下只有一个消费者活跃它挂掉以后另一个接管兼顾了顺序和高可用键共享则把相同 key 的消息路由到同一个 consumer同时允许多个 consumer 并行处理不同 key 的数据。我在生产环境里用得最多的是共享订阅和键共享。很多业务消息本身不要求全局顺序只要求同一业务维度的顺序比如同一用户的操作要按时间顺序处理那键共享就是最合适的。选错订阅类型会直接影响消费吞吐和顺序语义这也是我在排查消费延迟时最先检查的项目。2.3 多租户与跨地域复制Pulsar 的多租户模型由三层组成tenant、namespace、topic。你可以把 tenant 理解成租户边界通常是企业内的一个业务线或者外部客户namespace 是 tenant 下面的文件夹用来做数据隔离和策略配置topic 则是实际的消息通道。这样做的好处是不同业务团队的数据从一开始就在逻辑上隔离可以设置不同的消息保留周期、存储配额和权限策略不需要物理拆集群。跨地域复制是 Pulsar 另一个我很喜欢的能力。它的实现原理并不复杂每个区域内的 Pulsar 集群维护一份本地写入日志同时异步把消息复制到远端集群这样一组集群之间共享同一个全局拓扑客户端从任意一个区域都能读写。更关键的是如果某个地域发生故障你可以把消费流量切到另一个地域而无需修改应用代码里的 topic 路径。这里要提醒一句跨地域复制不是免费的。复制会产生额外的网络流量和存储开销消息的最终一致性也需要业务端接受。如果业务要求强一致那么应用层还需要补充补偿机制。这些都是议程里可能深入讨论的内容提前理解概念后听 session 会更有代入感。2.4 存储计算分离收益与代价存储计算分离听起来很美好但实际收益和代价是对等的。收益方面最明显的是扩缩容灵活。broker 可以根据连接数和吞吐动态调整bookie 可以按存储容量和数据冗余独立扩容。另一个好处是写入路径的隔离broker 和 bookie 处理各自的性能瓶颈不会像单体型消息系统那样一个磁盘抖动把整个集群拖垮。代价也很直接。首先是组件变多Pulsar 集群至少需要 broker、bookie、zookeeper新版本里也可能是元数据服务各一套部署和维护的复杂度明显上升。其次是网络要求更高broker 和 bookie 之间的毫秒级延迟直接决定端到端的写入性能如果只是随手搭在虚拟机里、网络质量很差性能表现会非常难看。所以我通常的建议是如果只是小规模试用单机 standalone 模式就行如果真要上生产一定要把网络规划、监控告警、容量评估都纳入考量。开发者的 cesh 看起来简单但背后的运维挑战会在规模上来之后暴露。3. 从议程看消息中间件的创新实践方向3.1 云原生部署从裸机到 K8s 的一步一脚印现在新项目部署 Pulsar很多团队直接往 Kubernetes 上放。这本身也符合云原生的大趋势但落地过程并不像“helm install”这么简单。Pulsar 在 K8s 上的关键矛盾在于 bookie 是有状态组件需要稳定的存储。如果你用的是本地 SSD那就要配置好节点亲和性和局部存储调度避免 Pod 漂移导致数据副本丢失如果用的是云盘的网络存储又要接受额外的延迟开销。broker 相对幸运因为它是无状态的可以通过 HPA 或者手动扩容快速调整副本数。比较稳妥的做法是把 broker 部署成 Deployment把 bookie 部署成 StatefulSet再配合 PDBPodDisruptionBudget防止节点维护时同时下线多个副本。我在自己的环境里测试过K8s 上跑 Pulsar 的难点其实不在启动而在升级和扩缩容后的稳定性。比如 broker 滚动升级期间客户端长连接会重连如果没有优雅排空机制瞬间的重连风暴就会导致集群短暂波动。这些内容在议程里应该有专门 session值得团队里负责部署运维的人仔细看。3.2 轻量计算Pulsar Functions 与消息流处理Pulsar 内置的 Functions 模块经常被低估。很多人觉得流处理就应该上 Flink 或者 Spark Streaming但其实很多简单任务用不上那么重的引擎。Pulsar Functions 允许直接在消息到达时做一些轻量的 ETL、数据过滤、字段转换甚至简单的窗口聚合而且支持 Java、Python、Go 多种语言。我实践下来Functions 最适用的场景是“单条消息级别的处理”。比如从原始日志中提取关键字段、把消息根据内容路由到不同 topic、调用外部接口做数据富化。这类任务如果丢给 Flink要起一个作业、配置 checkpoint 和资源管理反而显得笨重。Functions 的优势是它和 Pulsar 的 topic 天然集成发布、订阅、容错都是现成的开发成本几乎可以忽略。当然Functions 也有边界。跨事件的状态管理、复杂窗口计算、事件时间语义这些应该交给真正的流计算引擎。这个取舍逻辑我认为是议程里“创新实践”非常典型的一个样本不是在追求技术炫技而是在合适的层级解决合适的问题。3.3 生态兼容Kafka 协议适配的取舍Pulsar 提供 Kafka 协议兼容层意思是你可以用原有 Kafka 客户端连 Pulsar不改代码就能迁一部分流量。这个能力对存量系统很有吸引力尤其是一时半会儿迁不完的团队。但兼容层不等于完全相同。我在测试中发现Kafka 协议兼容对绝大多数基本生产消费接口都能正常工作但一些高级特性比如事务、消费组内某些管理操作、Kafka Streams 的深层依赖可能就覆盖不到。如果业务重度使用 Kafka 的高级特性还是要认真评估。另外就算协议兼容性能表现也不会和原生协议完全一致毕竟底层是不同的实现。我会把兼容层定位成“迁移的桥头堡”而不是“永久的生产依赖”。先通过兼容层把流量切到 Pulsar等验证稳定后再逐步改造客户端为原生协议这样整体的迁移风险会小很多。这个循序渐进的方式也正好呼应了开发者日里“实践”的题眼。3.4 稳定性与可观测性治理消息中间件的可观测性经常被忽视直到出问题才想起来。我的经验是消息系统的监控至少要从四个维度入手吞吐、延迟、积压、存储。Pulsar 暴露的指标很丰富重点要盯的是 broker 的入站出站速率、每个订阅的 backlog、bookie 的磁盘使用和读写延迟。议程里会有治理方向的内容我也比较认可。因为消息中间件一旦出问题影响往往是大面积的。比如某个 topic 的消费者处理能力下降backlog 就会快速增长如果超过存储配额甚至会影响同一 namespace 下面其他 topic 的消息写入。这种“同一个租户的故障扩散”比起单体应用的宕机更隐蔽必须有自上而下的配额管理和监控告警。建议团队在刚接入 Pulsar 的时候就定义好关键指标和告警阈值不要等到积压已经涨上来才去逐个排查。4. 实操体验把 Pulsar 跑起来并完成一次收发4.1 本地快速启动想快速体验 Pulsar最省事的方式是跑 standalone 模式。我用 Docker 的次数最多因为不用管本地 Java 和 Python 的依赖。docker run -it \ -p 6650:6650 \ -p 8080:8080 \ apachepulsar/pulsar:3.3.0 \ bin/pulsar standalone启动完你会看到 broker 和 bookie 在一个进程里都起来了。然后用命令行创建 topic 并发送一条消息# 进到容器里执行 docker exec -it container_id /bin/bash bin/pulsar-admin topics create persistent://public/default/demo-topic bin/pulsar-client produce persistent://public/default/demo-topic \ --messages hello-pulsar看到生产成功之后再消费试试bin/pulsar-client consume persistent://public/default/demo-topic \ --subscription-name my-subscription \ --num-messages 0如果--num-messages 0是持续消费模式你可以 CtrlC 终止。第一次跑通这个流程Pulsar 的基本机制就心里有数了。4.2 生产消费者代码演示命令行跑通以后建议再写一段客户端代码因为绝大多数生产场景都会用 Client SDK。下面是我常用的 Java 客户端示例基于 Pulsar 3.xPulsarClient client PulsarClient.builder() .serviceUrl(pulsar://localhost:6650) .build(); ProducerString producer client.newProducer(Schema.STRING) .topic(persistent://public/default/demo-topic) .enableBatching(true) .batchingMaxMessages(100) .batchingMaxBytes(4096) .batchingMaxPublishDelay(10, TimeUnit.MILLISECONDS) .create(); producer.send(user-1 login event); producer.close();消费端同样很直接ConsumerString consumer client.newConsumer(Schema.STRING) .topic(persistent://public/default/demo-topic) .subscriptionName(my-subscription) .subscriptionType(SubscriptionType.Shared) .receiverQueueSize(1000) .subscribe(); while (true) { MessageString msg consumer.receive(5, TimeUnit.SECONDS); if (msg null) { continue; } try { System.out.println(received: msg.getValue()); consumer.acknowledge(msg); } catch (Exception e) { consumer.negativeAcknowledge(msg); } }这段代码里有几个点值得留意。生产端开启 batching 能显著提升吞吐但会增加小量延迟适合对实时性没那么敏感的场景。消费端的receiverQueueSize决定客户端本地缓存多少条消息设得大一些可以让消费更平滑但也要防止应用崩溃时消息在本地队列里丢失所以需要配合业务需要设置。4.3 几个值得调的参数新手最容易犯的错误是直接使用默认配置然后到线上才发现吞吐或延迟不达标。我把自己常调的几个参数列出来供参考managedLedgerDefaultMarkDeleteRate控制消费确认标记删除的节奏默认值可能会让积压确认不及时调高可以加快存储空间释放。brokerDeduplicationEnabled开启生产消息去重可以防止客户端重试导致的重复写入但要接受一定的内存开销。tcpNoDelay默认关闭开启后能降低小消息的传输延迟。客户端receiverQueueSize这个不是 broker 参数但非常重要。太小会降低消费吞吐太大容易堆积内存建议根据消息大小和 JVM 堆内存综合评估。调参不是盲目抄网上配置而是要在压测过程中观察响应。参数调完记得用 Pulsar 自带的管理工具和监控指标对比改动前后的差异不然很容易“调了个寂寞”。5. 常见问题与排查技巧实录5.1 消费者 lag 持续上涨这是我在生产环境遇到过最多的场景。消费者 lag 上涨第一反应是“处理太慢”但实际原因可能有好几种。我见过最典型的三个原因第一订阅类型选错。如果业务量不小但用了独占订阅那就只有一个消费者在工作其余副本全闲着吞吐自然上不去。换成共享订阅或键共享把并发打开 lag 很快就能消化。第二单条消息处理过慢。比如消费者逻辑里调用了下游 HTTP 接口下游服务响应极慢整个消费链路也被拖住了。这种情况靠调队列参数没用要从业务处理链路优化。第三接收队列太小导致等待频繁。如果消费端receiverQueueSize10而消息到达速度很快消费者大部分时间都在等待网络请求而不是在处理消息。把队列调大以后吞吐立刻改善。排查 lag 最直接的方式是看 Pulsar 监控里的 subscription backlog 指标再结合消费者日志判断是哪个环节瓶颈。5.2 连接被频繁断开连接不稳定看着像网络问题很多时候其实是参数没配对。Pulsar client 和 broker 之间有 keepalive 机制如果客户端这边很忙长时间没有发送数据就容易被 broker 判定为无效连接然后断开。特别是跨网络访问的时候防火墙还会对长时间空闲的 TCP 连接自动清理。解决思路有三个方向一是调整客户端的 keepAlive 参数比如keepAliveIntervalInSeconds不要设得太长二是检查 broker 和客户端之间的防火墙规则确认不会清理空闲连接三是如果使用了负载均衡器要确保四层负载的 idle timeout 不小于心跳间隔。我曾经排查过一个问题应用每两小时固定报一次“connection reset”后来发现是防火墙的 idle timeout 是 60 分钟而客户端 keepalive 默认周期更长最后调整 keepalive 到 30 秒就恢复了。遇到这种问题不要急着怀疑 Pulsar先把网络链路的时间参数全部过一遍。5.3 磁盘占用快速增长Pulsar 的存储增长来源有两个正常写入的消息以及未被确认消费的 backlog。很多人以为停止生产后磁盘占用就会稳下来如果消费端没 ack消息数据会一直被 BookKeeper 保留磁盘只会越涨越高。出现这种情况先看是不是某个订阅没有活跃消费者维护游标。比如一个临时调试用的订阅运行一段时间后就不再消费但它的订阅游标一直留在某个老位置整个 backlog 都不会被清理。此时可以直接删除该订阅或者在 namespace 层面设置合理的保留策略和 backlog 策略。常见的做法是把消息保留时间设成业务可以接受的最大值比如 12 小时同时配置brokerDeleteInactiveTopic自动清理不活跃 topic。对于 Long-term 数据备份需求可以把老数据转移到分层存储层而不是在 hot 层无限堆积。5.4 几点建议最后的运维建议来自我的实际经验。第一刚接触 Pulsar 不要一上来就搭复杂集群先用 standalone 模式把 API、命令行、消息流程跑通建立手感。第二准备上生产前做一次容量估算至少算清楚峰值写入速率、平均消息大小、保留时长和副本数用这些参数推导出 bookie 的磁盘需求。第三提前设计好监控告警和配额管理让异常在早期就暴露出来。另外Pulsar 这类系统对网络质量很敏感。开发环境里可能不会暴露问题一旦跨机房、跨公有云链路出现抖动抖动就会被放大到消息延迟和写入失败上。如果有条件尽量做一些网络演练确保故障发生时运维有预案。我个人在使用 Pulsar 的过程中最大的感受是它的分层架构和多租户模型确实是为大规模和复杂组织准备的但这也意味着它需要更多的设计意识和运维纪律。这次 Pulsar Developer Day 的议程安排几乎每一个方向都在朝着“如何让这套系统更容易被驾驭”在走。如果你也有正在落地的消息中间件实践建议把议程里关于云原生部署和性能调优的部分重点看一遍哪怕只是从中找到一条能改进的思路都值回票价了。