ARTICLE DETAIL

资讯详情

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

Pulsar Key_Shared 模式消息不消费?从存算分离架构到哈希区间分配排查实践

Pulsar Key_Shared 模式消息不消费?从存算分离架构到哈希区间分配排查实践 COSCon‘25 的议程表单里Pulsar Developer Day 的日程一发布我第一时间就点进去翻了一遍。作为一个日常跟消息中间件打交道的人我对这种官方发布的行程向来不抱太高的期待但这次围绕 Pulsar、以“消息中间件创新实践”为主题的专场确实有几个点值得单独拿出来聊聊。这篇文章不准备给你复述议程条目而是换一种打开方式从议程透露出的热度聊聊 Pulsar 到底解决了什么问题、它和传统队列/流平台的分野在哪里再把最近社区里被反复讨论的“pulsar 的 key_shared 模式不消费”这个坑完整拆一遍。无论你是在调研选型还是已经把它跑在生产环境里这篇内容应该都能给你一点参考。1. 议程亮点速览Pulsar Developer Day 值得关注的几个方向1.1 一个消息中间件项目为什么值得单独办一个 Developer Day说实话消息队列这类基础组件在开源圈子里一直是“用的人多、愿意深聊的人少”。大部分用户停留在“发消息、收消息、最多看看 backlog”的层面能长期投入做贡献、做源码级分享的人是少数。Pulsar 能单独开一个 Developer Day至少说明两件事一是它的社区活跃度已经到了一定体量二是它确实有足够多的新鲜话题值得用一整天来展开。对于参会者来说这种专场的价值在于“信息密度”。平常你翻文档、看博客学到的是孤立的知识点而一个 Developer Day 的议程通常会把某个话题从原理讲到落地、再从落地讲回原理一天听下来你会对整个系统的运行逻辑有一个完整的认知。尤其对于正在做技术选型的团队这种集中输入比你自己花两周读源码高效得多。1.2 从议程结构来看今年重点围绕哪几类话题虽然不能替组委会把每个议题都剧透一遍但从已经公布的议程方向看Pulsar 相关的技术分享大致绕不开这几条线架构演进类存储与计算分离的实践细节、BookKeeper 的调优、Broker 无状态化运维、元数据服务的性能瓶颈。生产落地类跨地域复制、多集群容灾、大规模多租户管理、消费堆积治理。这类是用户最想听的因为线上踩过的坑才最有含金量。生态集成类Pulsar 与 Flink/Spark 的集成、Pulsar Functions 的适用边界、Kafka API 兼容层的真实表现。运维治理类容量规划、限流配额、负载均衡策略、事件驱动架构中的可靠性设计。我的观察是每次 Pulsar 相关活动里“踩坑复盘”型演讲的出场率都挺高。原因也简单Pulsar 的配置项多、机制复杂光靠官方文档很难覆盖所有组合场景而社区里那些“我把线上搞挂了后来查明白了”的分享恰恰是价值密度最高的内容。1.3 我的看法这类活动最大的价值是“踩坑复盘”我参加过好几次类似的开源专场最深的感受是听别人讲他怎么把系统架构搭起来远不如听别人讲他怎么把系统排查清楚。前者往往是结果导向你看到的是一个“已经稳定运行”的系统后者则是过程导向你能看到一个人面对异常指标时是怎么逐步缩小范围、最终定位根因的。这种思维方式很难从手册里学到只能在真实的故障场景中沉淀。所以如果你打算去现场我建议你把“某个模式在某些版本下表现异常”这类话题列为重点关注对象。后面我会花大篇幅讲的 key_shared 不消费问题就是一个非常典型的案例。2. 消息中间件的创新风向Pulsar 凭什么成为焦点2.1 存算分离不是噱头它改变了扩容和故障恢复的玩法Pulsar 和传统消息队列最底层的差异是它把计算和存储彻底拆开了。Broker 不保存数据只负责接收消息、分发消息和记录消费位点真正落盘的是独立的 BookKeeper 集群。你可以把这个架构理解成超市的收银台和仓库收银台Broker只管结账商品都堆在仓库BookKeeper里。收银台不够的时候加收银台仓库不够的时候加仓库两者互不干扰。这带来的直接好处是扩容优雅。Kafka 的 topic 和 partition 是绑在 broker 上的想扩容就要把分区分批迁移过去期间还有流量压力和数据一致性风险。Pulsar 的 topic 则被切成一个个 ledger 分布在 BookKeeper 上对 broker 来说数据存储位置是透明的。你要加 broker直接加一台机器让负载均衡器分配一部分 topic 过来就行不需要做数据搬迁。不过要提醒一句存算分离不是银弹。分离之后网络开销和元数据服务的压力会变大ZooKeeper 或者 etcd 这类协调组件反而更容易成为瓶颈。很多团队把 Pulsar 搬上线后遇到的第一个问题不是消息发不出去而是元数据操作超时。所以评估 Pulsar 时别只盯着“架构先进”这四个字还要看你有没有能力维护这套更复杂的依赖体系。2.2 多租户和订阅模型Pulsar 最容易被忽略的两张王牌Pulsar 的设计里有两个概念很容易被低估多租户Multi-tenancy和订阅模型。多租户体现在租户tenant/ 命名空间namespace/ 主题topic三层结构上。你可以在一个集群里给不同部门划分不同的命名空间每个命名空间独立设置存储配额、消息保留策略、消费权限。对中大型团队来说这意味着不用为每个业务单独搭一套集群节省的成本非常可观。订阅模型则是 Pulsar 区分于普通队列的核心。它提供了四种订阅类型订阅类型行为特征典型场景Exclusive一个订阅里只能有一个消费者消息严格有序顺序处理、单点消费Shared多个消费者轮流拿消息谁拿谁处理吞吐最高无状态、可乱序的并发处理Failover一个主消费者处理备用消费者只在主节点故障时接管主备容灾Key_Shared按消息 key 做哈希分发同 key 的消息固定到同一个消费者需要同 key 有序、又要并行处理Key_Shared 是这里面的“折中派”它不像 Exclusive 那样只有一个消费者也不像 Shared 那样完全不考虑顺序。它给相同 key 的消息固定分配消费者从而保证这些消息能被有序处理同时不同 key 的消息可以分发到不同的消费者手里去做并行消费。2.3 从 Kafka 迁到 Pulsar到底是换了赛道还是添了复杂度这几年确实有不少团队从 Kafka 迁到 Pulsar理由无非是这几点支撑更大规模的 topic 数、需要原生的跨地域复制、想要分层存储、希望用 Functions 简化流处理链路。这些需求 Kafka 不是不能做但往往要拼装多个组件才能实现Pulsar 把它们做进了内核里。但我也要给跃跃欲试的团队泼一盆冷水Pulsar 的整体复杂度明显高于 Kafka。它不是一个单一二进制就能跑起来的系统你要同时维护 Broker、BookKeeper 集群还有元数据服务。如果你所在团队没有专职的中间件运维能力我建议先从托管服务开始等真正理解它的运行机制后再考虑自建。选型这件事“能解决问题”只是及格线“你能养得活它”才是真正的门槛。3. key_shared 分区不消费的坑从一次线上事故说起3.1 现象消息只进不出消费者却一切正常最近社区里被反复提到的“pulsar 的 key_shared 模式不消费”问题我一开始是当个传闻看的直到自己线上也踩到了一次才发现这个坑比想象中隐蔽得多。当时的现象是一个业务用 Key_Shared 订阅消费某个 topic某天开始 backlog 异常上涨但所有消费者进程都是活的心跳正常日志里没有任何报错。用管理工具查看订阅状态消费者列表也在看起来一切正常。可消息就是卡在 topic 里一条都拉不出去。这种“没有任何明显异常、但功能失效”的问题最难受因为你连从哪里下手都不知道。3.2 一步步排查问题出在 hash range 的分配上我把当时的排查链路完整复盘一遍这比直接给结论更有用。第一步先确认基础设施没有故障。查看 Broker 的 CPU、内存、网络指标确认 BookKeeper 的写入延迟正常。这一步排除了底层故障导致的消息堆积。第二步用pulsar-admin topics stats和pulsar-admin topics stats-internal查看订阅详情。重点看两个地方当前订阅的 consumer 列表是否完整、每个 consumer 的 availablePermits 是否有余量。结果发现某个消费者虽然出现在列表里但它的 availablePermits 一直是 0而且它的位置始终没有变化。第三步对比 Shared 模式下的表现。同一个 topic 改用 Shared 订阅临时测试消息立刻就能被消费。这一下就把范围缩小到了 Key_Shared 模式特有的分发机制上。第四步引入 Key_Shared 的实现原理。Key_Shared 模式下Broker 会把整个哈希区间0 到 65535划分成若干段分给各个消费者。消息的 key 通过哈希函数映射到某个点这个点落在哪个区间段消息就交给对应的消费者处理。如果某个区间段没有消费者认领或者某个消费者注册了却没拿到任何区间段那么哈希到这个区间的消息就永远不会被消费。第五步确认根因。当时那个线上环境消费者数量发生了几次动态变化触发过几次 rebalance。在旧的消费者被剔除后Broker 为剩余消费者重新切分哈希区间时没有给新注册上来的消费者划分新的 range。结果就是有一个消费者“占了坑但没分到地”对应的那部分 key 的消息从此再也没有被拉取。3.3 为什么这个坑这么难发现因为这个 bug 的触发条件很隐蔽消费者进程健康、网络没问题、rebalance 也完成了从外部指标看一切正常。唯一的异常是“某个消费者的哈希区间为空”而这类信息在默认监控里根本看不到。另外不同版本的 Pulsar 对 Key_Shared 的分发策略实现差异很大。早期的实现比较简单消费者数量一变区间切分就容易出现边界问题。后来官方引入了更灵活的哈希区间策略比如自动扩展区间、更新哈希区间等功能但这些策略有的默认没开有的在旧版本上存在行为不一致这又给问题多添了一层不确定性。还有一个关键点Key_Shared 不能用累计确认cumulative ack。因为它一个消费者会负责多个 key 的消息如果对之前所有消息做一次性确认会把其他 key 还没处理完的消息也确认掉导致数据丢失。很多踩坑案例其实是 ack 方式用错了把 Key_Shared 当成 Shared 来写消费逻辑结果消息反复被投递、最终堆积。3.4 临时的救命解法和长久的根治方案如果你也遇到了类似的问题先别急着重构按下面这套优先级来操作。临时措施是这样重启出现问题的消费者进程强制 Broker 重新分配哈希区间。大多数情况下消费者重新注册后能触发一次完整的区间重算消息会恢复消费。如果重启没用就把这个订阅删掉让消费端重新订阅。注意这一步有风险删除订阅会丢失当前消费位点重新订阅时默认从最早或最新开始必须和业务方确认后再操作。在问题解决前控制消费者数量的动态变化不要频繁扩容缩容。每次变化都是一次 rebalance而 rebalance 正是触发这个 bug 的高危时机。长效方案就是这样优先升级到当前大分支的最新补丁版本。这个 bug 在社区 issue 区有多个反馈后续版本对 Key_Shared 的哈希区间分配做了不少修复升级是成本最低的解决方式。在消费端显式配置 KeySharedPolicy不要依赖默认行为。核心是打开自动更新哈希区间的功能让消费者数量变化后能自动触发区间重算。参考的 Java 客户端写法大致如下KeySharedPolicy keySharedPolicy KeySharedPolicy .stickyHashRange() .autoUpdateHashRange(true); ConsumerString consumer client.newConsumer() .topic(persistent://tenant/namespace/topic) .subscriptionName(subscription-name) .subscriptionType(SubscriptionType.Key_Shared) .keySharedPolicy(keySharedPolicy) .subscribe();注意不同客户端的版本 API 会有差异这个写法只展示思路落地时以你用的客户端版本对应的文档为准。3.5 把 Key_Shared 当成一种有纪律的模式来用经历过这次事故后我对 Key_Shared 的态度变了很多。它确实是一个好模式但使用它的前提是你理解它的运行规则并且有相应的监控手段。简单说三条纪律第一不建议为了单纯追求并发就把消费者数量拉得很高Key_Shared 的并发能力取决于 key 的分桶粒度不是消费者越多越好第二消费端必须使用单条确认不能偷懒做累计确认第三上线前必须验证消费者数量动态变化下的消费连续性这比压测峰值吞吐重要得多。4. 把 Pulsar 用稳的几条实战经验4.1 版本选型不要追新也不要停在太老的版本上处理 Pulsar 这类快速演进的项目版本策略其实很讲究。追新版本风险大因为大版本刚发布时往往带着一些社区还没踩完的雷但停在太老的版本更危险因为你积累的“稳定性经验”可能恰恰是建立在旧版本 bug 之上的。我个人的实践是选一个已经发布半年以上、有稳定补丁版本的线比如 3.0 之后的某个中期版本。升级时先升客户端、再升 Broker每一步都用灰度 topic 验证。千万别一口气把整个集群升级完那样出了问题连回滚都来不及。4.2 消费端参数这些细节决定了线上稳定性消费端的参数配置看着简单实际很容易踩坑。我列几个关键点ack 超时时间Key_Shared 模式下如果业务处理逻辑有瓶颈消息会在 ack 超时后反复投递把堆积问题放大。ack 超时不要设得太短至少要给业务处理留出足够余量。批量接收 vs 单条处理很多团队为了提升吞吐开了批量接收但每条消息后端都要查库或者调接口批量接收只会在客户端堆积一堆未确认消息反而加大管理压力。建议根据业务处理的耗时来定处理时间越长的业务越不要追求大批量。多个 topic 共用订阅Key_Shared 叠加多 topic 订阅时哈希区间是按 topic 独立分配的不要指望所有 topic 的消息都均匀落在消费者上需要单独验证。rebalance 触发频率消费端频繁加机器、减机器都会触发哈希区间重新分配。Interval 太短会让客户端反复断开重连线上现象就是“消费时断时续”。4.3 监控别只盯吞吐要看消费延迟和区间分配Pulsar 的监控指标很多但大多数人只盯吞吐量和 backlog。吞吐量高不等于健康backlog 上涨也不一定代表 broker 有问题可能是某个消费者卡住了。我建议至少把这几类指标加进告警消费延迟从消息生产到被消费确认的端到端时间这是最直观的业务指标。订阅 backlog按订阅维度的堆积量比只看 topic 维度的 backlog 更细。消费者可用许可availablePermits如果某个消费者的 permits 长期为 0说明它可能已经被“遗忘”了。哈希区间数量能让消费者自己上报负责的 hash range 数量这个指标在 Key_Shared 模式下尤其重要。如果你已经有 Prometheus可以在 Exporter 层把 Pulsar 的消费组指标接进来然后针对这些指标配置独立的告警规则。别等业务方反馈“消息怎么不动了”才知道出问题了。4.4 压测不只是压峰值更要压故障场景很多团队做消息中间件压测的时候只会测“每秒能吞多少条”这远远不够。消息中间件的稳定性不光体现在吞吐上更体现在异常恢复能力上。我的压测清单里通常包含这些场景杀掉一个消费者观察剩余消费者是否能在预期时间内接替它的哈希区间。杀掉一个 Broker观察 producer 和 consumer 的连接重连时间。往 topic 里写一个大消息比如超过 1MB观察分片和 Key_Shared 的组合行为。填充一批 key 分布极不均匀的消息对比业务处理时间和消费堆积曲线。这些场景都能暴露一些“配置正确但行为异常”的问题。像 Key_Shared 的区间分配 bug很多时候就是在“杀掉一个消费者再启动一个消费者”这个动作之后才暴露出来的。5. 写在最后参会与上手建议5.1 在现场怎么听收获才最大如果你手头已经在用 Pulsar这次 COSCon‘25 现场的 Pulsar Developer Day我的建议是带着问题去听不要带着笔记本去抄。把你自己集群的 Broker 版本、客户端版本、Key_Shared 配置方式准备好遇到相关演讲直接拿这些信息去交流。基础组件的问题往往不是“它坏没坏”而是“你这个版本在这个场景下到底有没有限制”这种信息在网上搜要搜半天在会场问一个维护者可能一分钟就得到答案。5.2 如果你还没用过 Pulsar该从哪里开始没有生产环境压力的情况下建议先从本地单机模式或者托管服务开始用很小的 demo 把四种订阅模式都过一遍。重点不是跑通而是观察不同模式下消费者增删时的行为差异。尤其是 Key_Shared 这个模式看完我这篇文章你会知道它的坑在哪但只有自己动手复现一遍才能真正理解哈希区间分配这件事。我在生产环境用 Pulsar 这几年的最大感受是它把选择权给了你但也要求你必须明白自己在做什么。像 Key_Shared 这种模式文档上只有几行字可一个配置组合不对就是一夜堆积。写这篇东西就是希望你在踩到类似的“玄学问题”时能早一步想到 hash range早一点定位到根因。
返回列表