ARTICLE DETAIL

资讯详情

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

Pulsar架构解析:存算分离与Key_Shared订阅模式

Pulsar架构解析:存算分离与Key_Shared订阅模式 1. 开场这个周末去开源集市找点“硬货”周末的 COSCon25 开源集市Apache Pulsar 的展台就在那等着你。如果你正好在会场里转悠想找一个既能聊实时消息、又能聊存储架构的摊位Pulsar 那个展位值得多停一会儿。不光是领贴纸、换周边更重要的是能直接跟项目的 committer 和 Contributor 面对面聊技术——这种机会平时在线上社区里可不太容易碰到。Pulsar 是什么先给还没接触过的朋友一句话说清它是一个云原生时代的分布式消息和流数据平台从 Yahoo 内部孵化出来后来捐给了 Apache 基金会。跟 Kafka 这类老牌消息队列相比Pulsar 最大的不同是“存储和计算分离”的架构这也让它天生更适合做多租户、跨地域复制和大规模分层存储。对正在折腾实时数据管道、微服务异步解耦、或者想统一消息和流两种场景的团队来说Pulsar 是个很值得认真评估的选项。这篇文章不打算做成活动流水账我主要想借 “Pulsar 参展 COSCon25” 这个由头把 Pulsar 的几个核心技术点拆开讲清楚它的存储架构强在哪四种订阅模式到底该怎么选以及一个最近社区里讨论得很热的真实问题——“key_shared 模式不消费”的 bug。这段排查经验是我在实际环境里踩过的坑也应该是很多 Pulsar 使用者迟早会撞上的问题提前知道能帮你省不少时间。2. 先搞清楚 Pulsar 到底解决了什么问题2.1 存算分离架构Broker 和 BookKeeper 的分工很多刚开始接触 Pulsar 的人会问Pulsar 和 Kafka 到底有什么本质区别我一般用一个类比来解释Kafka 像是把“收银台”和“仓库”绑在一起开分店每个分店自己管自己的货Pulsar 则把所有“收银台”集中在前台后面只有一个巨大的中央仓库由专门的库管团队BookKeeper来负责分拣和存取。具体到架构上Pulsar 的 Broker 层只负责消息的路由、权限管理、订阅状态跟踪这些相对轻量的逻辑处理而真正的数据存储落在 Apache BookKeeper 上。Broker 是无状态的你可以随时加机器横向扩容只要把新节点注册进集群流量就能自动分摊过去。BookKeeper 的数据节点Bookie按 Segment段来存储消息每个 Topic 的数据被切成一堆小 Segment分散存储在多台 Bookie 上并且每个 Segment 默认写多副本通常是 3 副本。这意味着某一台 Bookie 挂了它的副本会立刻顶上来数据不丢、服务不断。这个架构带来一个很实在的好处Broker 不需要本地磁盘存消息。Kafka 的分区数据是要落盘的Broker 挂了换新机器得先把数据重新同步回来迁移成本高得很。Pulsar 的 Broker 挂了新 Broker 直接接上数据还在 BookKeeper 那边躺着几乎是无感切换。对运维而言这个差异在故障演练的时候体会最深。2.2 Topic 可以很大、很多且支持多租户和分层存储另一个 Pulsar 让我觉得非常顺手的地方是 Topic 数量可以撑到百万级。早期的消息队列一旦 Topic 多了整个集群性能都会跟着崩原因在于每个 Topic 的元数据和索引都要占资源。Pulsar 在应用层把 Topic 抽象成了逻辑概念实际物理上是一个“Ledger 上的 Segment 集合”Broker 即使管理上万个 Topic自身压力也远小于传统方案。多租户方面Pulsar 用了三层命名空间结构Tenant租户— Namespace命名空间— Topic。比如一个 Topic 的完整名字是persistent://my-tenant/prod-namespace/my-topic。你可以在每个租户级别做隔离、配额、速率控制几个业务团队共用一套集群互不干扰。对中大型公司而言这直接省掉了“每个团队单独搭一套 Kafka”的运维成本。分层存储Tiered Storage也是我很看重的一项功能。默认情况下数据按保留策略在 BookKeeper 里滚动删除但只要你配了一个 S3 兼容的对象存储Pulsar 就会自动把老的 Segment 数据搬到 S3 里。读的时候如果消费者想回溯旧数据Broker 会从对象存储拉回来。这样你就不必为了“偶尔回溯一下三个月前的数据”而无限扩大 BookKeeper 集群的磁盘容量成本控制非常直观。2.3 从 Kafka 迁到 Pulsar 的体验观察我身边有不少团队做过 Kafka 到 Pulsar 的迁移反馈基本集中在几个点一是 Pulsar 的消息消费模型更灵活同一个 Topic 可以同时跑“队列式消费”和“流式发布订阅”二是管理端做得比较清晰pulsar-admin 命令行工具、REST API 和 UI 控制台都能操作三是社区对生态连接器比如 Pulsar IO Connector的支持已经相当丰富Source 和 Sink 一接就能用。当然不是说 Pulsar 就面面俱优。它的架构让每个环节都更分布式这意味着部署和运维的组件变多了至少要部署 ZooKeeper或者 etcd新版已经支持元数据服务抽象、BookKeeper、Broker 三套东西。小团队如果只有一两台机器想轻量跑动一个功能完整的 Pulsar 集群会比单机 Kafka 费劲不少。所以它的优势场景更偏向有一定规模、真正需要弹性扩容和稳定性的生产环境。3. 消费模型的根本差异订阅模式选对了事半功倍3.1 Pulsar 的四种订阅模式对比在 Pulsar 里订阅Subscription是消费行为的核心抽象。同一个 Topic不同的订阅之间是隔离的各自维护各自的消费进度。Pulsar 支持四种订阅模式订阅模式消费语义适用场景Exclusive独占一个订阅只允许一个消费者其他消费者连接时直接报错严格有序、单消费者场景比如全局事件表、严格有序的状态变更流Failover灾备多个消费者里只有一个 Active其他作为 StandbyActive 挂了自动切换需要顺序保证且要求一定可用性的场景比如金融交易事件Shared共享消息随机分发给多个消费者各处理各的不保证顺序吞吐量大、对顺序不敏感的任务型队列比如异步通知、批量任务分发Key_Shared按键共享相同 key 的消息只会路由到同一个消费者不同 key 可以分散到多个消费者既需要按 key 保证顺序又想并行扩容的场景比如用户维度的实时聚合、订单状态机流转3.2 Key_Shared 到底应该怎么用我把 Key_Shared 单独拿出来讲因为它在“顺序”和“并行”之间找到了一个很好的平衡点。举例来说你在处理一笔订单的实时数据流订单 ID10086 的“创建”“支付”“完成”三条消息如果落到 Shared 模式下很可能被两个消费者抢着处理最后出现“支付先于创建”这种乱序而用 Key_Shared只要我们把订单 ID 作为消息的 key这三条消息就一定会进同一个消费者顺序自然保住了。与此同时不同的订单 ID 仍然能分给不同消费者处理吞吐量不会因为单消费者而受限。在代码上用 Java Client 指定 Key_Shared 订阅很容易ConsumerString consumer client.newConsumer() .topic(persistent://tenant/ns/my-topic) .subscriptionName(my-sub) .subscriptionType(SubscriptionType.Key_Shared) .subscribe();有一点必须提醒Key_Shared 模式下消息的分发是“尽力按 key 哈希落到消费者”但它在设计上并不保证一个消费者独占一批 key 之后立即把该 key 的后续消息都交给它。也就是说如果你既要严格顺序又要重构某些聚合状态最好还是用 Exclusive 订阅或做补偿机制。生产上跑 Key_Shared还需要开启消费者端KeySharedPolicy并合理配置哈希分片数量默认的 2 个分片在小规模消费组下可能不够均匀。3.3 踩过的订阅模式选型坑以前接手过一个项目业务方想用 Key_Shared 同时实现“按用户 ID 有序”和“多个消费者并行”。听起来很完美跑了一阵子发现消费组里的某个消费者负载明显高于其他节点。查了半天发现是消息 key 的基数太小——比如 key 只取了业务类型下单、支付、退款而业务类型总共就 3 种哈希来哈希去就落在几个固定消费者上。这不算 bug但很容易被误判成 Pulsar 的分发策略有问题。后来我把 key 设计成“业务类型_用户ID”组合分布瞬间均匀了。像这类负载不均的问题如果想在测试阶段就发现可以先用pulsar-admin topic stats看一眼每个消费者的msgRateIn指标曲线能直接暴露出分发不均。别省这一步。4. “key_shared 不消费”问题实录现象、排查与解决4.1 事故发生时的现场这个周末开源集市的群里正好有人在聊“key_shared 模式不消费的 bug”我一下子就想起了几个月前在生产环境踩过的同款情况。当时我们的服务大概是这样的一个订单系统用 Pulsar Key_Shared 订阅消费订单变更事件一共起了 8 个消费者跑了一段时间后突然监控报警消费组的 Lag积压暴涨。但是诡异的是你去查消费者的状态每个消费者都显示“已连接”“正常”订阅仍然存在ranges 也还在就是不出消息。关键信息我再理顺一遍不是消费者崩溃不是订阅不存在也不是权限问题。broker 端没有任何异常日志Producer 那边发送消息的速率完全正常消息也确实写进了 Topic繁荣的消息就堆在那但就是没有消费者拉到数据。这种“连接正常但不消费”的状态比直接报错难查太多了。4.2 排查路径从客户端到 Broker 再到元数据我当时的排查顺序是这样的每一步都有明确的排除目的第一步先看客户端日志。在 Pulsar Client 里消费事件会记录在 debug 级别包括 Broker 分配的消费队列和 receiver queue 的 fill 情况。把日志级别临时调到 DEBUG发现消费者端根本没有收到新的MessageReceived事件。这里就说明问题大概率出在“Broker 没把消息匹配给这个订阅”上下而不是消费回调卡死。第二步看pulsar-admin的统计信息。执行bin/pulsar-admin topics stats persistent://tenant/ns/orders-events重点看subscriptions里的my-sub对应的msgOutCounter和msgRateOut。当时我看到的msgOutCounter完全不动而msgBacklog一直在涨这进一步把问题锁定在 Broker 该订阅的分发逻辑上。什么连接、心跳、TCP 层面的问题都可以先排除掉了。第三步查订阅模型和 Topic 策略。Key_Shared 模式有一个很容易被人忽略的约束它要求 Topic 上的所有订阅都必须是“有序消费类型”或者至少要让 Key_Shared 能够正常处理。如果同一个 Topic 上既有 Shared 订阅又有 Exclusive 订阅某些老版本 Broker 对 Key_Shared 消费者的路由表处理会异常。当时我很怀疑是某个排查用的调试脚本在同一个 Topic 上建了另一个普通 Shared 订阅从而干扰了 Key_Shared 的订阅状态。第四步也是最关键的看元数据里的 consumer 分配情况。当时我们团队先在一个测试 Topic 上复现了这个问题不用业务侧代码直接用pulsar-perf consume加上 Key_Shared 订阅去消费结果同样没消息。这就把问题缩小到了 Broker 源码层面的路由逻辑。查了相关 issue 和源码终于明白在较旧的 Pulsar 版本里Key_Shared 的消息分发依赖 broker 为每个 Key_Shared 订阅维护“hash range 到 consumer”的映射关系。一旦消费者做 rebalance 或有新的消费者加入、退出导致的 hash range 更新不及时Broker 会认为所有 range 都已经被“冻结”或没有可用的 consumer 而停止分配新消息。表现正是consumer 连接正常但消息不出。4.3 修复方案与规避手段最终我们升级了 Pulsar Broker 到修复该问题的版本这里也建议读者在生产环境长期关注 Apache Pulsar Release Notes 中关于 Key_Shared 的 bugfix 条目。如果短期没法升级有几种临时的规避方式可以用在订阅不变的情况下重启消费者进程强制触发一次 rebalance。把消费者的数量临时调整为 1再用后台任务业务补偿把积压消费完必要时再逐步增加消费者。如果积压时间敏感直接新开一个 Shared 订阅做兜底消费把积压先追平同时注意业务侧接受短时的乱序或者在下游做 key 维度的窗口排序。另外要检查客户端的版本和 Broker 版本是不是相差太大新客户端连老 Broker 的 Key_Shared 订阅在某些协议版本组合下兼容性并不可靠。这个 bug 给我的教训很深凡是依赖 Key_Shared 这种“按需有序”模型的业务必须配套一套“乱序兜底”的应急方案比如在下游加一个 keyed 合并缓冲层至少不能让消息彻底卡死在那里。消息系统再可靠也要对极端 bug 留一手。4.4 给你的 Key_Shared 使用建议清单作为一个被坑过的人我现在总结了一套 Key_Shared 上生产前的检查清单确认业务的顺序语义真的需要 key 粒度。如果你接受全 Topic 乱序直接用 Shared别多此一举。确认消息 key 的基数足够大避免少数几个 key 占满消费者。不使用同一个 Topic 混合不同订阅模型尤其不要混用 Shared 和 Key_Shared。所有消费者使用相同版本 Client并尽量与 Broker 保持同 minor 版本。为订阅配置合理的maxUnackedMessagesPerConsumer否则消息不 ack 会把消费卡死。持续关注 Pulsar 社区关于 Key_Shared 的 issue因为这个功能在多个版本里都有过不同程度的 bug。5. 从 COSCon 开源集市聊到生产实践里的 Pulsar 细节5.1 活动现场可以聊什么如果你周末真的去 COSCon25 开源集市了建议别只领贴纸就走。Pulsar 展台通常会有以下几样值得深度交流的东西一是项目核心开发者对架构演进方向的最新解读比如 Pulsar 3.x 里对元数据服务、负载均衡、Broker 自动故障转移的改进。平时在 Release Notes 里看到的文字远不如现场听作者讲一遍“当时为什么要这么设计”来得通透。二是关于 Pulsar 生态工具链的演示比如 Pulsar Functions轻量流处理、Pulsar IO Connector、Pulsar SQL用 Presto 查 Topic 数据这些组件能让你在不用引入一套 Flink 或 Spark 的情况下就做不少实时计算和查询的事儿。三是最重要的——现场“诉苦”环节。开源社区最宝贵的资源其实是那些踩过坑的真实用户。你可以拿着自己的配置、监控截图、问题栈直接跟 committer 聊。我上一次在开源大会的 Pulsar 展台就靠着一张“消费者连接正常但不消费”的截图让开发者帮我定位到了配置上的问题——那种效率远高于自己在 issue 列表里翻半天。5.2 生产环境里容易被忽视的 Pulsar 配置细节借此机会我再把几个生产环境容易被大家忽视的 Pulsar 参数拎出来聊一聊。第一个是确认ackTimeout和ackTimeoutRedeliveryBackoff的关系。默认情况下 ack timeout 如果设置得太短而业务处理偶发超过这个时间就会造成消息被重复投递消费者被迫做大量幂等逻辑。如果你对消息时延没那么敏感建议把 ackTimeout 调到业务处理 P99 时长的 5 倍以上或者干脆关闭自动 ack timeout改用negativeAckRedeliveryBackoff配合死信重试。第二个是managedLedgerDefaultEnsembleSize、managedLedgerDefaultWriteQuorum、managedLedgerDefaultAckQuorum这三个 BookKeeper 参数。默认值通常是 2/2/2 或 3/3/3决定了 Segment 的副本数。如果集群里只有 3 台 Bookie你配了 3 个副本坏一台机器后可用性会受影响。副本数也直接决定磁盘用量和写入吞吐先想清楚业务允许丢失多少数据再定这几个值。第三个是消费者端的receiverQueueSize。它相当于是你的消费窗口大小。队列太小消费吞吐上不去在 Pulsar 上面表现为“平均每秒钟只能消费几千条”队列太大一条慢消息会把整个消费窗口卡住几秒下游 RT 瞬间冲高。这个值一般建议设在 1000 左右高吞吐场景可以试 5000但一定要配合maxUnackedMessagesPerConsumer一起调。第四个是关于“Backlog 管理”。不少人只关注消费 Lag没想过积压消息如果迟迟不消费BookKeeper 磁盘会一直不释放。Pulsar 的retention策略和ttl策略需要配合使用先设 TTL让过期消息自动清理再设 retention保证已经 ack 的消息还能被回溯一段窗口。如果都不设消息永远删不掉磁盘迟早撑爆。5.3 开源集市上的 Pulsar 周边以及如何带一个问题回家往年 COSCon 开源集市的惯例是每个项目展台都会有印章收集或者小任务集齐几个可以换周边。Pulsar 的周边经常是飞行棋、贴纸、帆布袋一类质感在同行的开源项目里算不错的。但对我来说比周边更有价值的是带着“现场验证一个真实猜想”的心态去聊。比如你可以在去之前先在本地环境里起一个 Pulsar Standalone 集群# 下载并解压 Pulsar 后直接启动单机模式 bin/pulsar standalone --num-bookies 1然后建一个 Topic用 Java 或 Python 客户端分别测试 Exclusive、Shared、Key_Shared 三种订阅模式下的行为。如果你想提前踩一下上面说的那种坑可以把 Broker 版本固定在某个比较旧的 2.x 版本然后频繁增删消费者观察 hash range 的变化。现场把这些发现跟 Pulsar 的开发者一聊你会发现自己对消息系统的理解立刻高了一个level。6. 常见问题速查现场最容易听到的 Pulsar 提问结合活动现场和社区群里听到的高频提问我做了一个“速查表”大家遇到类似问题可以直接对着查现象常见原因解决方案Producer 发送报Topic not found没有启用自动创建 Topic 或权限不足Broker 配置allowAutoTopicCreationtrue赋予对应 role 的生产消费权限消费有延迟但 CPU 不高ack 太慢单消费者内积压了太多 unacked 消息增大receiverQueueSize调大maxUnackedMessagesPerConsumer消费重复消息明显变多ackTimeout 太短消息未处理完被重复投递调大 ackTimeout配置重试退避而非立即重投单 Topic 吞吐无法突破分区数太少创建分区 Topic确认分区 Producer 是否按 key 正确路由某个消费者永远没消息可能是 key 哈希全部被其他消费者占用或处于 pending rebalance检查msgRateIn和consumerName触发一次 rebalance 观察bookie磁盘使用率持续增长没有 TTL 或 retention 策略旧消息不被清理设置ttl和retention配合pulsar-admin topics delete清理无效 Topic在 COSCon 展台现场另一个问得非常多的问题是“Pulsar ZooKeeper 挂了怎么办”。这个问题的本质是Pulsar 只用 ZooKeeper 管理元数据并不承载消息数据的读写。ZooKeeper 短暂不可用Broker 和 Bookie 之间的现有连接还在工作消息数据读写不受影响只是元数据变更比如 Topic 创建、订阅状态变更会暂时不可用。所以生产环境把 ZooKeeper 组件单独部署、监控告警配好一般不会出现整个消息集群不可用的极端情况。7. 聊聊我对 Pulsar 在开源社区发展的一点感受在 COSCon 这类开源集市上我观察到一个很明显的变化前几年摆摊的大多还是 Web 框架类项目大家问的是“这个东西能干什么”现在越来越多像 Pulsar 这样的基础软件项目出现在集市上来看的人直接拿着一张集群拓扑图或者一堆异常监控截图来问“这里该怎么调”。这个变化背后是基础架构层面的人才密度在上升也是大厂开源项目逐渐走向生产可用、走向成熟的信号。如果你周末到了现场大概率会看到 Pulsar 展台前围着的人分成两类一类是刚入行、眼里发光的学生另一个是带着生产环境难题来的老工程师。这两类人之间的对话往往才是开源集市上最精彩的环节。中间的技术布道者会反复解释同一件事——为什么消息系统需要存储计算分离为什么流和队列可以在一个平台里统一为什么 Pulsar 愿意把自己最核心的 BookKeeper 存储层单独开源出来。这种“愿意剖开自己、让更多人理解底层设计”的氛围是我觉得国内开源社区最珍贵的东西。又想起那个 key_shared 的 bug。其实任何一个能扛住大规模生产流量的开源项目它的成长路径里都写满了类似的 bugfix。能把它拿到台面上在 issue 里、PR 里、社区活动里被反复讨论本身就是项目生命力的体现。这也是为什么我更推荐大家去现场和开发者聊聊——你踩到的每一个坑都有可能是推进下一次版本迭代的一块砖。周末如果有空来 COSCon25 开源集市逛逛吧。可能你站在那里跟 Pulsar 开发者聊了十分钟带回的不只是周边而是一整套关于自己系统架构的新思路。
返回列表