ARTICLE DETAIL

资讯详情

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

Pulsar消息中间件架构选型:对比Kafka/RocketMQ

Pulsar消息中间件架构选型:对比Kafka/RocketMQ 做消息中间件选型这些年我前前后后用过 RabbitMQ、RocketMQ、Kafka每个都在特定场景下扛过不少流量。但真正让我觉得“这套架构思路有点东西”的还得算 Apache Pulsar。Pulsar 给我的第一印象是“它不像传统 MQ更像一个自带存储引擎的流数据平台”。这篇就顺着架构决策这条线聊聊我在评估和落地 Pulsar 时的完整思考包括它解决了什么问题、对比其他 MQ 的核心优势、实操部署中必须注意的细节以及那些面试里高频出现的 MQ 问题幂等、重复消费、消息顺序在 Pulsar 里到底怎么解。如果你正处于“要不要从 Kafka 迁到 Pulsar”或者“新项目选型到底选哪个 MQ”的纠结阶段这篇应该能帮你把决策路径理清楚。我不会只讲概念会尽量把架构层面的取舍、压测参数、踩坑记录都摆出来方便你直接拿去参考。1. 为什么架构决策中我会认真考虑 Pulsar1.1 Pulsar 到底解决了什么痛点先说结论Pulsar 最大的贡献是把“消息中间件”从“无状态路由层 外部存储”的拼凑模式变成了“计算与存储分离”的原生架构。传统消息队列比如 RabbitMQ依赖 Erlang 虚拟机做路由和队列管理消息堆积能力有限Kafka 用分区日志存储依赖分区数量来扩展吞吐但分区多了会带来文件句柄、rebalance 成本、Broker 状态同步等一系列问题。Pulsar 从一开始就把这些拆开了。Pulsar 的核心组件分为两层Broker 层负责生产和消费请求的处理是无状态的BookKeeper 集群负责消息数据的持久化存储是有状态的。这两层可以独立扩容。Broker 不够了加 Broker存储不够了加 Bookie互不干扰。这个架构带来的直接好处有三个一是扩容时不需要搬迁数据Broker 是无状态的加机器就能分担流量二是存储和计算可以按需独立扩展不像 Kafka 那样每个 Broker 既要做网络 IO 又要管理本地磁盘三是多租户隔离做得更干净不同团队共用集群时各自的存储和流量策略可以独立控制。我在实际评估时的感受是如果团队业务处于快速扩张期消息量从每天几百万涨到几亿Kafka 的扩容成本会越来越高——每次加节点都可能触发分区 rebalance需要提前规划操作窗口有限。而 Pulsar 扩容 Broker 几乎对客户端无感Bookie 扩容也能做到在线完成。这个差异在架构决策表里是 Pulsar 非常重要的加分项。1.2 Pulsar 与 Kafka、RocketMQ 的选型对比很多人在做架构决策时会纠结到底选 Kafka 还是选 Pulsar我的建议是不要只看吞吐量数字要结合业务模型来判断。如果只是做日志采集、大数据管道Kafka 的生态更成熟配套的 Connect、Streams、Flink 集成更丰富团队上手也快。但如果业务中消息队列承担的是核心交易链路对消息不丢失、分区数扩展、多租户隔离要求较高Pulsar 的架构优势会更明显。具体对比可以从几个维度看。存储模型上Kafka 每个 Broker 使用本地磁盘存储分区副本分区数量和存储量受机器磁盘限制Pulsar 使用 BookKeeper 集群统一存储数据分散在多个 Bookie 上broker 不落盘。消息消费模式上Kafka 只有 partition 级别的 offset 管理要新增消费者组做重复消费就得重置 offsetPulsar 的 subscription 机制允许同一份数据被多种订阅模式消费而互不影响。我做过一个不算严谨但很直观的对比测试同样 3 个节点、机械硬盘环境下Kafka 单个分区写入吞吐大约能到 10MB/s 级别增加分区能提升整体吞吐Pulsar 单 Topic 的写入也能稳定到类似水平但因为存储层在 Bookie 上做了分段写入和缓存对磁盘 IO 的抖动容忍度更高。震动厉害的云磁盘环境中Pulsar 的延迟曲线比 Kafka 更平滑。RocketMQ 的对比更有意思。RocketMQ 也是存储与 Broker 同节点但它做了很好的顺序写优化性能不差事务消息支持也成熟。可它的分区扩展同样麻烦Topic 的读写队列数在消息发送后变更会有风险。Pulsar 因为存储和计算分离队列数partition 数量可以随时调整不会影响已有数据。团队技术栈如果已有 ZooKeeper 依赖选 Pulsar 需要额外接一个元数据中心Pulsar 这边叫metadata storeKafka 从 2.8 开始逐步去 ZooKeeper用 KRaft 模式RocketMQ 有自己的 NameServer。这一点也是架构决策中的隐性成本需要提前算进去。2. 核心架构与设计思路拆解2.1 Pulsar 的分层架构Broker 与 BookKeeper 各司其职Pulsar 的架构分层不是噱头它是理解 Pulsar 一切特性的钥匙。先说 Broker它的职责是接收客户端连接、解析协议、鉴权、管理 topic 的 ownership、把读写请求转发给对应的 Bookie。因为 broker 不存数据所以它能做到水平扩展非常平滑你可以在流量高峰前直接加 broker元数据服务会自动做负载均衡。BookKeeper 是 Pulsar 的存储引擎它采用segment分段存储机制一个 topic 的消息被拆成多个 ledger 段每个 ledger 包含若干 entryentry 按顺序追加到 Bookie 的 journal 文件和 log 文件里。Journal 是预写日志保证崩溃恢复时数据不丢log 文件负责最终落盘。这种分段设计让 Bookie 的 IO 模式变为顺序写、随机读和传统数据库的 B 树随机写完全不同机械磁盘上也能跑出不错的性能。这里有一个容易忽略的点Pulsar 的数据可靠性依赖 BookKeeper 的副本机制。每个 ledger 默认写入 2 个副本可通过配置修改并分散到不同机架的 Bookie 上。写入时客户端会等待至少 quorum 数量的 Bookie 确认后才返回成功。这个机制和 Kafka 的 ISR同步副本集合思路类似但实现是在存储层完成的Broker 无需关心副本同步细节。实际运维中Bookie 的磁盘配置直接影响性能。我踩过的一个坑是BookKeeper 的 journal 盘建议使用高性能 SSD而 log 盘可以用普通 HDD。如果混用同一块盘当数据量上来时journal 的写入延迟和 log 的刷盘会互相干扰导致写延迟波动变大。生产环境给 Bookie 至少配两块盘是底线。2.2 订阅模型与消费语义不止是发布订阅Pulsar 的订阅模式是它区别于其他 MQ 的大亮点。它支持三种订阅类型Exclusive排他订阅、Shared共享订阅、Failover灾备订阅。Exclusive 模式下一个 subscription 只允许一个消费者Shared 模式下消息按 round-robin 分发给多个消费者适合吞吐量大的场景Failover 模式下主消费者接收所有消息备用消费者在主消费者故障时接管保证消息处理的连续性。这个设计解决了一个业务痛点同一份消息数据不同团队如果要用不同的消费方式只需各自创建独立的 subscription互不干扰。举个例子订单系统产生的订单事件交易团队用 Shared 订阅做并行处理风控团队用 Failover 订阅保证严格顺序两者消费进度互不影响。在 Kafka 里你需要复制 topic 或通过 consumer group 管理 offset实现成本和心智负担都更高。消费语义上Pulsar 提供 at-least-once、at-most-once 和 effectively-once 三种级别。默认是 at-least-once配合 consumer 端实现幂等可以做到业务上的不重不漏如果要精确一次可以使用 Pulsar 的Message Deduplication功能和Transactional消息。实际开发中大多数业务做到 at-least-once 加幂等消费已经足够没必要为了“精确一次”牺牲吞吐。顺带提一个经常被问到的热词“前端点两次算是发两条消息吗”。这个问题的本质是业务请求重复提交但 MQ 能否自动去重。Pulsar 有生产者端去重功能开启后生产端传入相同的MessageId或SequenceIdBroker 会丢弃重复消息。但前端连点产生的两次请求大概率消息内容不同比如带了不同的时间戳所以最终还是得靠业务层幂等。这个我会在后面的“MQ 幂等问题”里详细展开。3. 实操落地从部署到生产的关键环节3.1 单机快速体验与集群部署要点如果你想快速体验 Pulsar可以直接用 Docker 跑单机版。命令很简单拉取apachepulsar/pulsar镜像映射 6650客户端端口和 8080HTTP API 端口即可。官方文档里有完整步骤这里不赘述。单机版适合做功能验证和学习但别拿它做性能测试因为单机没有体现分层架构优势。真正上生产建议至少部署 3 个 Broker、3 个 Bookie 和一个元数据中心可以用 ZooKeeper 或 etcd。这里有一个容易忽视的点Pulsar 的元数据中心不要和 BookKeeper 的元数据混在同一个 ZooKeeper 集群里用尽量用独立的 ZooKeeper 或者直接用 Pulsar 的metadata store配置独立的 etcd 集群避免相互影响。部署时还有一些细节需要提前规划Broker 的managedLedgerDefaultEnsembleSize默认 2和managedLedgerDefaultWriteQuorum默认 2决定了每个消息写入的副本数如果机器少或磁盘紧张可以降到 1但生产环境我建议保持 2 或 3。Bookie 的journalMaxLocations、dbStorageWriteCacheMaxSizeMb这些参数需要根据磁盘和内存调整。官方默认值偏保守直接跑没有任何问题但想榨干性能还是得调优。初次部署最容易出错的是网络端口配置。Pulsar 对外暴露的端口有 6650二进制协议、8080HTTP、8000Admin 接口需根据版本确认Bookie 之间还有 3181 端口用于 BookKeeper 通信。如果有防火墙或安全组记得把这些端口都打开。我曾经在一台云主机上部署客户端怎么都连不上排查半天发现是安全组没放行 6650这个基础问题最容易被忽略。3.2 生产者与消费者的使用姿势Pulsar 客户端的 API 设计很简洁Java 和 Python 都支持。生产者创建时推荐手动指定topic、producerName、sendTimeout并开启enableBatching批量发送来提升吞吐。这里强调一下sendTimeout的设置如果客户端与 broker 之间的网络有抖动过短的 timeout 会导致大量发送失败重试加重 broker 压力建议设置到 30 秒以上并配合retry配置。消费者使用上subscriptionName是核心概念同一个消费者组的实例必须使用相同的subscriptionName才能让服务端把消息分发给组内成员。receiverQueueSize参数控制消费者本地缓存的消息数量默认 1000如果业务处理每条消息耗时较长可以适当调小比如 100避免消息堆积在客户端内存里。对于 Shared 订阅的消费者每条消息处理完后必须调用acknowledge否则消息会在ackTimeout后重新投递。在实际开发中最容易出现的问题是消费端忘记开启ackTimeout或者设置不合理。Pulsar 的 ack 机制不是自动确认必须在业务处理成功后手动 ack。如果处理逻辑抛异常消息会被重新投递配合消息重试策略就形成了“业务重试”。这里要注意重试消息会进入replay队列如果业务本身有依赖外部系统比如第三方接口的场景需要设计好退避重试策略避免集中重试打垮下游。3.3 端到端压测我如何评估 Pulsar 的性能压测是架构决策中不可或缺的一环。我的压测思路分三步单 Topic 单分区吞吐、多 Topic 多分区扩展性、消息积压场景下的表现。单 Topic 压测时用官方提供的pulsar-perf工具命令类似./bin/pulsar-perf produce -r 10000 -n 100000 my-topic。实测下来3 个 Broker 3 个 Bookie 的集群单 Topic 写入吞吐约在 8 万到 12 万条/秒消息大小 1KB开启批量。这个数字仅供参考因为实际吞吐取决于磁盘性能、网络带宽和消息大小。多 Topic 扩展性测试更有意义创建 100 个 partition 的 Topic用 20 个生产者并发写入吞吐能线性增长到 40 万条/秒左右。注意Pulsar 的分区数不像 Kafka 那样受 Broker 数量限制理论上分区可以设置得很大但分区太多会占用 Bookie 的文件句柄和内存。建议根据预期流量来计算分区数单分区吞吐约 1 万条/秒如果预期 10 万条/秒设置 12~16 个分区比较稳妥。消息积压场景测试是很多人忽略的。我先用消费端 sleep 模拟慢消费让消息积压到几千万条然后恢复消费观察 Broker 和 Bookie 的负载。Pulsar 在积压场景下表现很好因为消息在 Bookie 上分段存储消费可以并发读取多个段而不会像 Kafka 那样出现单分区瓶颈。我在压测中把积压量拉到 5000 万条消费恢复后吞吐稳定在 5 万条/秒没有出现 OOM 或磁盘 IO 打满的情况。这一点对于秒杀类业务非常关键。4. 面试与团队协作中常被问到的细节4.1 MQ 幂等问题到底怎么解“MQ 幂等问题”是面试热点也是实际开发最容易踩坑的地方。先说结论消息中间件层面只能保证 at-least-once 投递做不到全局不重复所有“不重复”需求最终都归结到消费端的幂等设计。在 Pulsar 场景下消费端幂等有几个常用方案。第一种是唯一键去重在消息体里带上业务唯一 ID比如订单号、流水号消费者处理前先查 Redis 或数据库唯一索引如果已存在则直接返回。第二种是数据库乐观锁更新每条消息携带版本号更新时用update ... where version ?更新的行数为 0 时说明是重复消息。第三种是状态机校验消息里包含目标状态只有当前状态正确时才能推进比如订单状态从“待支付”到“已支付”重复的“已支付”消息不会影响当前状态。这里有个细节需要注意消费端实现幂等时要考虑事务边界。如果在处理消息过程中先改数据库再查 RedisRedis 写入失败会导致重复消息再次执行时误判为未处理。建议把“记录已处理状态”和“业务数据变更”放在同一个事务里比如通过数据库本地消息表先插入消息记录再更新业务表二者用同一事务保证原子性。Pulsar 的生产者端去重只能解决“生产者重复发送”的情况前提是客户端设置了消息去重编号SequenceId。但业务上的重复提交比如前端连点生成的是不同的 SequenceId中间件无法感知所以最终还是要在业务层做幂等。这也是为什么说“前端点两次”这个问题答案不能仅仅依赖 MQ。4.2 消息丢失与消息顺序两个绕不开的经典题消息丢失是一个系统性工程问题不是单一中间件能解决的。Pulsar 中消息丢失可能发生在三个环节生产者端、Broker 端、消费者端。生产者端丢消息常见原因是发送时没有开启确认或者使用了异步发送但没处理异常。Pulsar 的sendAsync回调里必须记录失败日志并做重试不能只调用不处理回调。Broker 端丢消息一般是副本数配置过少比如ensemble1或者 Bookie 磁盘故障。生产环境建议至少 2 副本并按机架感知配置避免一台物理机宕机导致数据全部丢失。消费者端丢消息通常是因为业务代码里先 ack 后处理或者处理完没 ack 但消费端重启。正确顺序是处理好业务逻辑、落库、最后 ack。消息顺序问题在 Pulsar 里比在 Kafka 里容易处理。Kafka 的顺序保证依赖分区内有序要求业务把同一 key 的消息投递到同一分区。Pulsar 中如果使用 Exclusive 或 Failover 订阅可以天然保证整个 Topic 的消息顺序如果使用 Shared 订阅则只保证单条消息的投递顺序不做全局有序。实际业务中订单状态变更这种需要严格有序的场景我用 Failover 订阅比较多而日志异步处理这种可以把顺序放开用 Shared 订阅提升吞吐。还有一个关键点是 key_shared 订阅模式。Pulsar 在 Shared 订阅基础上发展出Key_Shared模式相同 key 的消息会路由到同一个消费者且该消费者内按顺序处理。这个模式对于“同一个用户的操作必须顺序执行”的场景非常有用而且比 Exclusive 的吞吐高很多。我在一个支付回调系统中就是用 Key_Shared 订阅按userId做 key既保证了每个用户的操作顺序又能多消费者并行处理不同用户的事件。4.3 MQ 面试题高频点与 Pulsar 特色问题在面试中聊 MQ常规问题基本绕不开这几个如何保证消息不丢失、如何保证消息不重复消费、如何保证消息顺序、如何实现延迟消息、如何做消息积压处理。这些问题在 Pulsar 里都有比较清晰的答案路径。延迟消息功能 Pulsar 原生支持通过deliveryDelay设置消息的延迟投递时间。实现原理是 Broker 把延迟消息暂存在delayed delivery tracker中到达时间后才投递给消费者。这个机制比 RabbitMQ 的死信队列实现延迟更优雅也比 Kafka 依赖外部定时器方案更合理。Pulsar 的特色问题一般集中在几个点BookKeeper 的分段存储机制、Broker 无状态带来的扩容优势、分层架构下如何做流量隔离、事务消息实现原理。这些问题考察的不仅仅是 API 使用还有对底层存储和分布式系统的理解。准备时建议从“为什么 Pulsar 能做到这些”的角度去讲而不是背概念。有一个容易被问倒的细节是Pulsar 的 Topic 数和分区数到底怎么规划。我的实践经验是Topic 数尽量控制在千级别以内不要为了隔离业务创建上万个 Topic因为每个 Topic 都要在元数据中心和 BookKeeper 中维护元数据。分区数根据流量定但单分区吞吐超过 1 万条/秒时就要考虑增加分区。这里的“1 万”不是硬性指标具体和机器性能有关但作为一个预估基准够用了。5. 常见问题排查与避坑实录5.1 生产环境踩过的 5 个典型坑整理了一下这一年多在生产环境遇到的高频问题每个都是真实案例希望能帮大家绕开这些坑。坑一Bookie 磁盘写满后集群进入只读模式。表现是生产消息发送超时消费者也拉不到新消息。原因是 BookKeeper 检测到磁盘使用率超过阈值后会把 ledger 标记为只读。排查和处理方式是提前配置磁盘阈值告警并在部署时用独立的 mount 点挂载 Bookie 数据目录如果已经触发只读需要立即清理无用的 ledger 或者扩容 Bookie不能简单重启恢复。坑二消费者 group 的 subscription 名冲突。多个团队共用集群时如果不同业务用了相同的subscriptionName会造成消息被不相关的消费者抢走。事故现场是 Order 服务突然消费不到任何消息而日志发现另一个团队的消费者在连续 ack。排查后发现是两边工程里硬编码了同一个 subscription 名。规范做法是给每个业务线设置命名空间namespace并约束 subscription 的命名规则格式建议为{团队}-{业务}-{场景}。坑三批量生产者的 maxPendingMessages 参数太小。高并发写入时生产者发送请求堆积在客户端报ProducerBusy异常。解决办法是调大maxPendingMessages默认 1000但这个值和内存占用成正比不能无限调大。一般 2~4 核的客户端设置为 5000~10000 比较合适同时配合批量发送能有效降低 pending 量。坑四云服务器上 Bookie 的 journal 和 data 目录用了同一个云盘。表现是写入吞吐波动大偶尔出现高延迟毛刺。这是因为云盘的 IOPS 有限journal 的同步刷盘会阻塞 log 的写入。解法是拆分云盘journal 用高 IOPS 的 ESSDlog 用普通云盘并且做 raid 或 LVM 条带化提升吞吐。坑五启用消息去重后生产者的 SequenceId 没有保证单调递增。如果 SequenceId 不是递增的Broker 可能在乱序时误判重复消息。官方要求同一个 producer 的 SequenceId 必须严格递增如果你在业务里用随机字符串做去重 key会出现消息被错误丢弃的严重问题。解决方式是只在需要精确去重的场景开启 dedup并且对 SequenceId 做严格管理。5.2 监控告警与数据迁移实操生产使用 Pulsar 必须有完善的监控。官方提供 Prometheus 指标暴露接口默认在:8080/metrics可以接入 Grafana 展示。重点监控指标包括Broker 的brk_ml_NumTopicBuckets、brk_storage_ml_EntriesSizeBookie 的bookie_server_BOOKIE_READ_BYTES、bookie_server_BOOKIE_WRITE_BYTES以及磁盘使用率、JVM GC 时间、客户端连接数。数据迁移是另一个高优实操场景。我们当时把一个 v2.8 的集群升级到 v2.11过程比较顺利关键步骤是先升级所有 Bookie再升级 Broker。Pulsar 官方保证滚动升级的兼容性但建议参考版本的 release note。跨大版本升级时先在小集群上做演练确认pulsar-admin broker update和bookkeeper upgrade都通过后再动生产。如果要跨集群迁移数据有几种方案Geo-replication同步跨集群复制适合多活场景pulsar-admin topics peek-messages能做小批量的消息查看离线方案可以用pulsar-client consume落文件再用pulsar-client produce重新灌入。我们这个量级的迁移用了“双写 消费重放”的方式先让生产端同时写旧集群和新集群然后切换消费到新集群最后验证数据一致性后下线旧集群。整个过程没有停机但需要做好消息顺序和时间戳校验。5.3 从 Kafka 迁移到 Pulsar 的成本与收益架构决策最后都会落到成本和收益的计算上。从 Kafka 迁到 Pulsar成本主要在几个方面一是集群从“存算一体”改成“存算分离”运维要同时维护 Broker 和 Bookie 两套组件二是客户端语言绑定要多适配 Pulsar 的 API三是原来基于 Kafka 生态的组件如 Kafka Connect 的 source/sink需要重新找替代方案。收益方面我们最直接的感受是扩容效率。之前在 Kafka 上业务流量翻倍时我会花一整个晚上的窗口期去加分区、重启集群、观察 rebalance全程不敢掉以轻心。用 Pulsar 后加 Broker 就是拉起一个实例注册到元数据中心流量自动均衡。Bookie 扩容也是一样新加的节点自动参与新 ledger 的写入不需要手动迁移数据。另一个收益是成本优化。因为 Broker 无状态CPU 密集的 Pod 和存储密集的 Pod 可以分开调度云上资源利用率更高。我们原先 Kafka 节点为了存储容量必须买高配机器换 Pulsar 后可以让 Broker 用小的计算型实例Bookie 用大容量型实例整体成本大约下降了 20%~30%。当然迁移也有阵痛。团队的运维经验需要重新积累比如 BookKeeper 的参数调优、ledger 回收策略设置、以及 Pulsar 的生态工具链如 Pulsar Manager要花时间熟悉。如果团队对 Kafka 已经很熟悉且业务体量不大完全可以继续用 Kafka没必要为了追新而迁移。架构决策的核心始终是匹配业务需求而不是追逐技术热度。6. 我的一线经验和最后的建议做了这么多次 MQ 选型和落地我的体会是不要被某个中间件的宣传词冲昏头脑要把自己的业务场景写成一个需求清单逐项对比后再做决定。Pulsar 的存算分离架构确实先进但如果你的场景对延迟极其敏感要求毫秒级端到端Pulsar 的跨网络读写可能会让你犹豫反过来如果你经常要面对消息积压、分区扩容、多团队共用集群Pulsar 的收益非常明显。最后分享一个小技巧评估任何消息中间件前先在压力测试中模拟“最坏场景”——比如 Bookie 宕机一台、Broker 宕机一台、网络抖动 10 分钟看整体表现。我因为提前演练过失控场景在真实故障来临时才没有手忙脚乱。消息中间件是分布式系统中最核心的链路之一选型时多花的一周时间会在之后的一年里以十倍的价值回馈给你。
返回列表