ARTICLE DETAIL

资讯详情

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

Pulsar KeyShared不消费与NACK关闭问题全解析

Pulsar KeyShared不消费与NACK关闭问题全解析 这周末的COSCon‘25开源集市我在Pulsar的展位前站了大半天。按说开源集市的氛围应该是聊周边、聊社区、聊八卦结果被问到最多的问题反倒异常具体有人打开手机里的监控截图问我“KeyShared模式下消费者不消费到底是不是个bug”还有人在群里蹲了很久开口就问“客户端的negativeAcknowledge到底怎么关掉”。这两个问题我在线上被问过无数遍这次在线下又亲眼看着它们反复出现干脆回来整理成一篇长文把这两个点彻底讲透。这篇文章适合正在维护Pulsar生产集群、或者刚刚开始用KeyShared订阅和NACK机制的同学。文章不会停留在“怎么配”会更偏重“为什么会出现这些问题”“哪些是bug、哪些是设计如此”“按什么顺序排查最不容易漏”。如果你只是想把配置抄走第2.3节和第3.3节可以直接用如果你想知道社区里那些“不消费”“重复投递”的讨论到底在说什么建议从头读完。1. COSCon‘25 开源集市上Pulsar 展位一天聊了些什么1.1 集市问询和线上 issue 最大的区别线上的问题描述往往经过思考整理提问者会贴出版本号、日志、复现步骤。集市上的问题不太一样大家是路过看到Pulsar的Logo顺手掏出手机就开问。有人是生产环境已经炸了趁着周末来线下找答案有人是还在选型阶段想借机会实际问问踩坑情况还有人的问题其实和工作没有直接关系就是想搞清楚Pulsar官网里那些看起来很抽象的订阅模式到底是什么意思。一整天下来我发现一个规律大家的问题不在“这功能怎么用”而在“这功能为什么是这个表现”。比如KeyShared模式官方文档写得很清楚同一把key的消息只会固定发往同一个消费者。但很多人读完文档依然会在生产环境里踩坑因为文档不会告诉你“如果你所有消息都没设置key在KeyShared模式下它们会被当成同一个key来处理”更不会告诉你这种情况下你的消费者集群会退化成只有一个消费者在干活。这些信息只有在真实环境里被坑过的人才写得出来。1.2 为什么这两个问题总能被反复问起KeyShared和NACK这两个话题能成为集市上的高频问题原因其实很朴素它们都直接决定“消息到底还能不能正常消费”。KeyShared模式一旦出现某个消费者不干活业务侧看到的就是消息堆积、消费延迟、订单处理变慢。NACK一旦配置理解错了业务侧看到的就是消息在无限重投、重复消费、下游被打爆。这两个问题表面上看一个是订阅路由一个是确认机制好像没什么关联。但往深里看它们本质上都在说同一件事Pulsar的投递语义和单机消息队列差别很大很多东西不是看两眼配置就能猜对的。单机队列的“谁消费、失败了怎么办”都是简单直白的分布式环境下就变成了“key的hash范围怎么分、重投延迟怎么算、ack超时了谁来负责”。下面这两部分我会把这两条线分别展开。2. Key_shared 模式“不消费”先别急着怀疑 bug按这条链路一步步查2.1 KeyShared 的分发逻辑和大多数人以为的不一样先把KeyShared模式的底层逻辑讲清楚。Pulsar里有四种订阅类型Exclusive、Shared、Failover、KeyShared。Shared模式是所有消费者瓜分消息谁有空谁拿消息之间没有绑定关系。KeyShared则是在Shared的基础上加了一个约束消息按照key的哈希值被固定路由到某一个消费者同一个key的消息永远只进同一个消费者的接收队列。这个设计意图不难理解是为了在“多消费者并行消费”和“同一个key的消息不被打散”之间取得平衡。典型场景是订单事件同一笔订单的全部状态变更消息都带同一个订单ID作为key业务上必须按照顺序处理。如果走Shared模式同一笔订单的“创建”和“支付成功”两条消息可能被两个消费者同时拿到顺序完全失控。KeyShared就是为了解决这种场景而生的。但这里有个很多人没意识到的细节KeyShared模式下如果消息没有设置key那么这些消息会被Pulsar当成同一个key来处理。也就是说如果你的业务消息大量不带key又用了KeyShared订阅那这些消息全部会被路由到同一个消费者身上。这时候你看到的现象就是“一个消费者忙死其他消费者闲死”。这不是bug这是KeyShared设计下消息没有key时的一种兜底策略但它在生产环境里造成的观感和bug一模一样。用个生活化的类比KeyShared就像银行里按号分流的VIP窗口一个VIP客户同一个key只能被同一个窗口消费者接待。哪怕旁边窗口空着这个VIP客户也不能随便换人。如果你的VIP客户数量特别少又特别集中那某个窗口就会一直忙其他窗口看起来就像在摸鱼。2.2 “看起来像 bug”的四个典型症状和真实成因结合社区反馈和我自己的使用经历KeyShared模式下“不消费”的报障大概可以归成四类症状一新增消费者后消息仍然只往旧消费者那边走。这是被当成bug最多的一类。明明消费者扩容了但积压的消息就是不往新消费者身上分。真实成因是KeyShared的key分配是固定的一个key在某个时刻绑定在哪个消费者上就是哪个消费者不会因为集群里多了个消费者就自动把存量key重新分匀。新消费者能接到的只有后续新产生的、恰好被哈希到它身上的key。如果业务key就那十几个那新消费者大概率长时间拿不到消息。症状二某个消费者一直积压其他消费者空闲。这类问题十有八九出在key的分布上。比如业务方把“用户ID”当成key但某个头部用户的订单量是普通用户的几千倍那么这个超级用户ID对应的所有消息都会压在同一个消费者上其他消费者再空闲也帮不上忙。还有前面说的那种情况大量消息没有key全被当成同一个key塞给了一个消费者。症状三消费位点完全不动broker端的积压持续上涨。如果连“哪个key在哪个消费者上”都还没查先别急着怪KeyShared。这种整体不动的问题常见原因有三个消费者进程卡住了、消费者的receive队列被大消息占满导致新消息拿不到、客户端和broker之间的长连接断了但客户端没有及时感知。这三种情况跟KeyShared本身没什么关系但很多人习惯性先往KeyShared上怀疑。症状四消费者重启或断线重连之后消息不往重连后的消费者身上投递。这类现象和症状一类似都是sticky特性在起作用key在新消费者加入后并不会立即重新分配必须等原消费者离开订阅组broker才会把它的key范围转交给其他消费者。而且这个“转交”也不是瞬时的中间可能有一段窗口期。社区里关于这类问题的反馈非常多有一部分属于设计如此有一部分是真的跟broker版本实现有关。我说句实在话这四个症状里症状一、症状二、症状四本质上都是KeyShared的特性导致的不是逻辑错误。真正算得上bug的更多是broker在key范围重分配过程中出现的异常比如某个消费者退出后它的key范围没有正确转交、或者某些batch消息在重投时始终进不了正确的消费队列。这类问题社区确实有记录在apache/pulsar的GitHub仓库里搜“KeyShared”能看到大量类似讨论。2.3 一份可以直接照着做的排查清单如果线上真的遇到了KeyShared“不消费”我建议按下面这个顺序查每一步都有明确判断标准。第一步先确认“不消费”的是整个订阅还是单个消费者。用pulsar-admin命令看topic的统计信息pulsar-admin topics stats persistent://public/default/my-topic重点看输出里subscriptions部分的msgBacklog和consumers列表如果msgBacklog在持续增长所有消费者的blockedConsumerOnUnackedMsgs都是false说明消息在正常分发只是某个消费者的处理能力跟不上问题不在KeyShared在消费逻辑。如果msgBacklog在增长但只有一个消费者有msgBacklog对应的积压其他消费者的积压都是0那基本可以锁定是key路由不均。如果consumers列表里只有一部分消费者显示了connected状态先检查没连接的消费者为什么没连上再看版本。第二步检查消息key的分布情况。这一步是判断key路由是否均匀的关键。如果你有Pulsar生态里的工具就直接用没有的话可以用一个简单脚本统计最近一段时间里进入topic的消息key分布# 伪代码思路 consume messages from topic for each message: if message.hasKey(): count[message.getKey()] 1 else: count[__no_key__] 1如果__no_key__数量巨大或者某个key占了绝大多数那key路由不均衡的原因就找到了。这个统计在测试环境做一遍就行生产环境不需要长期跑。第三步检查消息是否有batch。如果你的producer开启了批量发送默认是开启的批大小一般10MB或1000条那么broker在KeyShared模式下可能把一整批消息作为整体路由给同一个消费者。这会导致某些本应分散到不同消费者的key因为处同一个batch里而被集体送到了一个消费者手上。这个问题在社区里被讨论过很多次如果你的业务消息都很小、又开启了高吞吐的batch配置建议在测试环境对比一下关掉batch前后的消费均衡情况。第四步查broker日志和客户端日志。在broker日志里搜KeyShared、reassign、unload等关键字在客户端日志里搜KeyShared、Reconnect。有一点很关键如果broker在key范围重分配过程中出现过异常日志里通常会留下线索。不用怕日志多把时间窗口对准“不消费”现象出现的时间段集中看那几分钟的日志就够。2.4 确认是 KeyShared 问题之后有哪几条退路如果排查完确认问题出在KeyShared本身就看业务场景选择对应的处理方案。退路一换订阅模式。如果业务上其实不要求同一个key的消息顺序处理直接换成Shared模式最简单。Shared模式没有key绑定的概念消费者之间有消息就会分出去负载最均匀。退路二改造消息key的基数。如果业务必须保留key顺序语义但又希望消费者分布更均匀可以在生产端对key做一层改写。比如原本只用“用户ID”当key可能导致超级用户把所有消息压在一个消费者上可以改成“用户ID 业务类型”或者“用户ID 分片后缀”让同一个用户的不同类型消息也能分流到不同的消费者。前提是你要能接受“同一用户不同业务类型的消息顺序不再严格保证”。退路三控制消费者数量变化的频率。KeyShared模式下频繁扩容缩容不会带来均衡收益反而可能引发key范围反复转移。如果一定要扩尽量一次性扩到位不要在短时间反复调整。退路四升级版本前先在测试环境复现。如果你怀疑是特定版本上的实现问题先把生产环境的版本组合broker版本客户端版本在测试环境完整搭一套用同样比例的key分布和消费行为去压测确认是不是版本问题再决定要不要升级。千万不要直接在生产环境升级客户端版本去做验证。3. 关闭客户端 negativeacknowledge先搞清楚要关的到底是哪个开关3.1 NACK 和 ackTimeout同一条投递链路的两套开关NACK这个话题我观察到一个普遍现象很多人说“我要关掉negativeacknowledge”但追问下去发现他们其实并不确定要关的是哪个开关。Pulsar的消息确认模型里跟NACK相关的至少有三套机制把它们放在一起看会清楚很多。机制触发方入口/API触发条件效果常见误解主动NACK业务代码consumer.negativeAcknowledge(messageId)业务代码显式调用消息被重新投递默认延迟1分钟以为NACK是自动发生的ackTimeout客户端consumerBuilder.ackTimeout(时长)消息在超时时间内没有ack客户端自动重新投递该消息以为关闭NACK就能关掉ackTimeout死信策略客户端brokerconsumerBuilder.deadLetterPolicy(...)消息重投次数超过阈值消息转入DLQ topic只配了NACK没配DLQ导致无限重投重点解释一下NACK和ackTimeout的关系。NACK是业务方主动告诉broker“这条消息我没处理成功你重新投递给我”。ackTimeout则是客户端自己干的活它会给每一条分发给消费者的消息启动一个计时器如果在超时时间内消费者没有确认那么客户端会自动把这条消息当成处理失败并触发重新投递。这两个东西很容易搞混是因为从结果上看它们都会导致“消息被重新投递”。但一个是业务层的主动行为一个是客户端层的兜底行为。很多人以为“我只要不调用negativeAcknowledge方法消息就不会重投”结果忘了自己配置了ackTimeout于是消息依然在无限重投。反过来也有人以为“我把ackTimeout关了NACK就没用了”结果业务代码里还留着主动NACK调用消息一样会重投。顺手说一句DLQ千万不要漏掉。如果你只关了NACK不做死信配置失败消息不会重投了但也一直留在消费位点那里不动后面的消息全被堵住消费卡死。这不是Pulsar的问题而是你没有给失败消息设计出口。3.2 为什么有人想关 NACK三个真实的业务动机“关掉NACK”这个需求本质上不是想禁用某一个API而是想改变失败消息的投递语义。三个最典型的业务动机是这样的。动机一重投导致的乱序问题。有一个真实的例子用户在处理订单事件时先收到“支付完成”事件又收到“订单创建”事件因为“订单创建”在第一次处理时失败被NACK重投了结果重投之后业务里出现了“创建一个已经支付过的订单”的情况。对于强顺序要求的业务NACK重投很可能放大了乱序窗口。动机二下游系统扛不住重复压力。有个做短信服务的朋友跟我说他们的消息处理依赖一个第三方短信接口这个接口一次只能承受有限的QPS。NACK重投对上游来说是自动的但下游接口不会理解“这是重投”于是每条失败消息的每次重投都会砸向接口把已经脆弱的服务彻底打垮。动机三希望失败快速走进死信流程。有些业务本身有幂等和重试机制消息消费失败只是偶发业务方希望失败消息不要无限重试而是快速进入死信队列由人工介入。这时候NACK的无限重投只会拖慢整个链条。这些动机都是合理的。但我要提醒一句关闭NACK本质上是把“失败后的自动重试”换成了“失败后必须由你自己兜底”如果业务侧没有准备好兜底方案关闭带来的问题比不关还大。3.3 真正“关闭”NACK 的配置组合以 Java 客户端为例真正要做到“关闭NACK相关行为”需要同时处理三件事不调用主动NACK、关闭ackTimeout、设置合理的重投延迟和死信策略。下面这段是Java客户端的示例配置PulsarClient client PulsarClient.builder() .serviceUrl(pulsar://broker.example.com:6650) .build(); ConsumerString consumer client.newConsumer(Schema.STRING) .topic(my-topic) .subscriptionName(my-subscription) .subscriptionType(SubscriptionType.Shared) // 关闭 ackTimeout 自动重投设为0 .ackTimeout(0, TimeUnit.SECONDS) // 如果不打算主动NACK可以把重投延迟调到一个很大的值避免意外触发 // 但实际生产里更建议保留合理延迟配合下方DLQ使用 .negativeAckRedeliveryDelay(60, TimeUnit.SECONDS) // 死信策略兜底重投超过3次进DLQ .deadLetterPolicy(DeadLetterPolicy.builder() .maxRedeliverCount(3) .deadLetterTopic(persistent://public/default/my-topic-dlq) .build()) .subscribe(); while (true) { MessageString msg consumer.receive(); try { // 业务处理 consumer.acknowledge(msg.getMessageId()); } catch (Exception e) { // 如果决定关闭NACK这里不要调用 consumer.negativeAcknowledge(msg.getMessageId()) // 建议让消息走DLQ或者打日志后人工处理 consumer.acknowledge(msg.getMessageId()); // 手动确认避免消费位点卡死 } }这段代码里发生了几件事逐一解释ackTimeout(0, TimeUnit.SECONDS)这是关闭客户端自动重投的关键配置。不设置这个字段也可以因为Java客户端默认不启用ackTimeout但显式写出来更容易让后来维护的人看懂意图。negativeAckRedeliveryDelay(60, TimeUnit.SECONDS)控制当业务代码确实调用了NACK时消息延迟多久重投。如果你真的想把NACK“关掉”更大的值是更安全的但也要注意值设太大意味着消息迟迟不确认可能卡住后续消息。这里用60秒是给最常见的配置方式。业务catch里不调用NACK而是手动ack。这一步很关键失败消息如果不确认也不NACK它就会一直留在未确认状态最终阻塞消费。手动ack之后消息从消费位点移除同时因为配置了DLQ消息在累计重投超过3次后会被自动转入死信topic。这样失败消息既不会无限重投也不会堵住后面的消息。Go和Python客户端的逻辑是一样的只是API名字略有差异。Go客户端是ConsumerOptions.NackRedeliveryDelayPython客户端是consumer.negative_acknowledge方法。核心思路不变不调用主动NACK、关闭ackTimeout、配置DLQ兜底。3.4 关闭之后最容易翻车的三个连锁反应配置好之后建议在测试环境观察一下这几个连锁反应。连锁一不手动ack也不NACK消费位点会长时间卡住。我见过一个真实案例业务方把NACK调用删了但忘了处理失败消息结果某一条失败消息一直占据消费位点后面所有消息都在等它超时。这个坑非常常见处理逻辑上一定要有“失败后手动ack或走DLQ”的出口。连锁二关闭NACK后消息失败是静默的。如果不打日志、不接报警失败消息进入DLQ之后你可能完全无感。建议在死信topic的消费者上接一个监控和报警这样失败消息进来时能及时感知。连锁三ackTimeout关闭后网络假死难以及时发现。开启ackTimeout时如果消费者进程和broker之间的连接假死消息会因超时被重新投递业务侧很快会看到重复消息并发现问题。关闭ackTimeout后消息会一直堆在消费者未确认区直到连接真正断开才发现。这个取舍需要业务方自己权衡。4. 从集市上的闲聊看客户端版本治理升级和验证两手抓4.1 很多人的“Pulsar bug”其实是客户端和 broker 版本错配集市上有个用户跟我说了一个很有意思的经历他们生产环境用的是Pulsar 3.0的broker但客户端的版本还停留在2.10结果KeyShared模式下出现了几种很诡异的消息不消费现象。后来他们统一把客户端升级到了与broker匹配的版本问题就消失了。虽然我没法把这个案例百分之百归因于版本错配但从社区经验来看客户端和broker版本差距过大会引入非常多不确定因素尤其是在KeyShared这类规则复杂的模式里。Pulsar的客户端和broker之间是存在版本兼容矩阵的。老的客户端连新的broker通常能工作但在某些新特性上会出现行为不一致新的客户端连老的broker可能直接遇到API不识别的问题。官方文档里给出了版本兼容性说明升级前一定要看一眼。我的建议很简单同一个大版本内尽量用较新的patch版本跨大版本升级前先看release notes不确认的地方最好在测试环境模拟一遍。很多线上问题其实是版本错配带来的而不是功能本身的问题。4.2 一套可复用的验证清单不管你是准备升级客户端、修改订阅模式还是调整NACK配置我强烈建议先在测试环境跑一遍下面这张验证清单。这也是我在集市上跟人聊的时候反复强调的“别拿生产环境当你的试验场”。验证项怎么做通过的判断标准KeyShared路由用固定key集合发送消息观察消费者接收情况同一key始终落在同一消费者新增消费者后存量key不迁移无key消息行为发送一批不带key的消息到KeyShared订阅确认这些消息的流向判断key分布是否符合预期NACK行为显式调用NACK观察重投延迟和重投次数重投延迟符合配置超过最大重投次数后消息进入DLQ关闭ackTimeout设置ackTimeout为0模拟消息处理超时观察消息是否不再被自动重投DLQ兜底构造失败消息触发重投上限消息进入DLQ topic消费位点没有被卡住版本兼容用生产版本组合跑一遍完整消息链路生产流量全链路无异常监控指标正常这张表不重要重要的是你在跑的时候要有真实的消息量和故障注入不能只发三条消息就宣布验证通过。拿一件故障注入工具或者手动模拟消费者崩溃、broker重启、消息处理抛异常这些场景把配置的真实行为在测试环境里彻底摸一遍比看十遍文档都有用。4.3 集市交流中一个比较深刻的观察最后说一个我这几天在现场的真实感受。开源集市这种场合大部分人来问问题的时候都是在生产环境已经出过事之后才想着补课。给我看KeyShared问题的用户说“之前只是照着教程配了没仔细想过消息key的分布”问NACK的人说“我们是上了线之后才看到无限重投的监控报警”。这其实反映了分布式消息队列使用中一个普遍的问题很多配置项不看说明书也能跑起来但只有出问题时才会明白它们到底在干什么。所以我的建议一直都是花点时间把自己的topic消费模型画出来标明订阅模式、消息key设计、失败处理策略和死信出口然后对照文档逐项检查。这个过程看起来很笨但它是发现潜在问题最高效的方式。我在实际查看KeyShared问题的时候还有一个很小但很有效的技巧处理“一个消费者积压”的问题时先看broker端topic stats里每个consumer的msgBacklog差值再去看消息key分布不要一开始就陷入代码逻辑。多数情况下几步就能定位到根因。这个技巧希望也能帮到你。
返回列表