拓冰建站拓冰建站
首页 / 资讯中心 / 正文

Pulsar消息重试与死信机制详解:原理、配置与生产避坑

1. 先解决“消费失败之后怎么办”这个根本问题做消息队列开发和运维的人迟早会在生产环境撞见这么一幕消费者把消息拉下来业务处理抛出异常消息被退回 Broker然后又被投递下来又异常循环往复。你要是没配置任何兜底策略这条消息就会像个赖着不走的钉子户卡死整个消费位点后面的消息全部积压。这段时间我一直在跟 Pulsar 消息重试与死信机制打交道从踩坑到理顺今天就把这套机制的底层逻辑、配置要点、排查经验一次讲透。这套东西解决的核心问题很简单消息处理失败之后系统应该怎么做。重试负责给“暂时失败”的消息第二次机会死信负责给“永远失败”的消息一个明确归宿。对于做支付回调、订单状态同步、外部接口对接的同学来说这套机制是必须吃透的。因为分布式环境下网络抖动、服务重启、依赖超时都是常态你不可能保证每条消息都一次处理成功所以“失败后怎么处理”就成了系统健壮性的分水岭。先说清楚一个误区很多人觉得消费失败就把消息重新塞回队列就行但其实这远远不够。如果你不做任何配置Pulsar 默认会反复投递同一条消息直到达到ackTimeout或negativeAckRedeliveryDelay的控制条件但这个过程可能是“无脑重试”既没有延迟策略也没有次数上限。结果就是一条处理不了的消息会反复打扰你的消费者甚至把 CPU 和日志磁盘打满。你真正需要的是有序的重试节奏以及一个“实在不行就让消息去死信队列”的终止条件。这套机制适合谁来学两类人一类是正在用 Pulsar 做业务开发的后端工程师另一类是维护 Pulsar 集群的中间件/SRE。前者需要知道怎么在客户端配置重试和死信策略后者需要理解 Broker 侧的工作机制才能在出问题时快速定位。今天的内容两种视角都会覆盖。2. 重试机制拆解ack、nack 和重试主题2.1 确认机制是重试的地基Pulsar 的消息确认ack机制是一切重试逻辑的基础。消费者处理完一条消息后显式调用acknowledge(messageId)Broker 才会认为这条消息被成功消费并把它从待确认状态里移除。反过来如果消费者没确认或者说主动调用了negativeAcknowledge(messageId)简称 nack那这条消息就会被重新投递。很多人刚接触时会把 ack 和 nack 搞混。这里有个非常关键的点ack 是通知 Broker“这条消息我搞定了”nack 是通知 Broker“这条消息我处理失败了你重新发给我”。nack 之后的消息不会立刻重投而是要等一段延迟时间这个延迟由negativeAckRedeliveryDelay控制。默认情况下Java Client 里这个值是 60 秒也就是说你 nack 了一条消息最快也要 60 秒后才会再次收到它。这个延迟参数是重试机制的节拍器。设置得太短比如 1 秒那遇到下游接口抖动时你的消费者会疯狂重试每秒打爆下游设置得太长比如 10 分钟那业务的恢复时间会被拉长用户感知到的延迟就很大。我一般建议从 5 到 30 秒起步根据下游服务的响应时长来调。如果下游接口 P99 是 800ms那 5 秒的重试间隔已经比较安全了。2.2 重试主题不是简单的“重新投递”Pulsar 的重试机制比 Kafka 的retry配置更进一层因为它引入了“重试主题”retry topic的概念。默认情况下你消费的是某个 topic比如persistent://public/default/orders下的消息而你配置了重试策略之后Pulsar 客户端会额外创建一个以-retry结尾的主题比如persistent://public/default/orders-retry。重试消息不是直接被塞回原主题而是先投递到重试主题中由客户端内部的消费者监听这个重试主题根据每条消息的时间戳判断“该不该重新投递了”。这背后用到了 Pulsar 的延迟消息投递能力——每条消息在重试主题中都有一个预期的投递时间时间到了才真正回到主逻辑中重新处理。这种设计的价值在于重试与主消费相互隔离不会因为重试流量影响主队列的消费顺序。而且你可以针对不同的异常类型设置不同的重试延迟。比如数据库死锁导致的失败等 10 秒再试大概率能成功外部 API 返回 500可能等 30 秒更合适。这些都可以通过重试主题和延迟级别实现。2.3 客户端重试配置实操我用 Java 客户端举例因为生产环境里 Java 用得最多。假设你有一个orders主题消费者代码如下ConsumerString consumer client.newConsumer(Schema.STRING) .topic(persistent://public/default/orders) .subscriptionName(order-sub) .subscriptionType(SubscriptionType.Shared) .negativeAckRedeliveryDelay(10, TimeUnit.SECONDS) .enableRetry(true) .subscribe();核心配置就两个negativeAckRedeliveryDelay控制 nack 后的重投延迟enableRetry(true)开启重试。但这里有个容易忽略的点如果只开enableRetry不配重试策略那么它默认使用的是MultiTopicConsumerImpl里的 default retry topic也就是在原始主题名后面加-retry后缀。如果你需要更精细的控制比如设置最大重试次数那就得自己定义DeadLetterPolicy这部分我会在第 3 节细讲。现在先记住重试的本质是“nack 延迟 重投”而重试主题只是实现延迟的手段。补充一个我在生产环境踩过的坑如果你开启了重试但消费端代码里没有调用negativeAcknowledge()而是直接不 ack、也不 nack那条消息会一直停留在消费队列里不会进入重试主题。最终它会等到客户端关闭或者ackTimeout超时触发 nack 行为后才会走重试。所以记得设置合理的ackTimeout否则消费者异常退出后消息会长时间卡在 unacked 状态。3. 死信机制给消息一个明确的终点3.1 死信的工作原理重试不是无限循环的总得有个“认输”的机制这就是死信Dead Letter存在的意义。当一条消息重试次数超过你设定的阈值Pulsar 会将它投递到一个专门的主题——死信主题DLQ默认命名是原始主题加-dlq后缀。死信机制不是 Pulsar 独有的Kafka 和 RocketMQ 都有类似设计。但 Pulsar 里实现方式有一点很贴心你可以在消费者级别配置死信策略而不需要额外创建一个消费者去监听 DLQ。当客户端判断一条消息的重试次数达到上限它就会自动将消息写入死信主题然后正常 ack 掉原始消息保证主流程不卡死。这个“自动”背后是有代价的Broker 不会主动帮你判断“消息失败了几次”是客户端在本地计数。这意味着如果你有多个消费者实例在消费同一个分区那么重试计数是分散在各实例内存里的。极端情况下一条消息可能被实例 A 重试 2 次、实例 B 重试 2 次合计 4 次但任何一个实例看到的都只有 2 次没达到阈值死信不会触发。这个问题在共享订阅模式下尤其明显后面排障部分我会专门讲。3.2 死信策略配置示例想在客户端启用死信最直接的方式是配置DeadLetterPolicy。下面是一个生产环境常用的配置DeadLetterPolicy dlqPolicy DeadLetterPolicy.builder() .maxRedeliverCount(3) .deadLetterTopic(persistent://public/default/orders-dlq) .build(); ConsumerString consumer client.newConsumer(Schema.STRING) .topic(persistent://public/default/orders) .subscriptionName(order-sub) .subscriptionType(SubscriptionType.Shared) .negativeAckRedeliveryDelay(10, TimeUnit.SECONDS) .enableRetry(true) .deadLetterPolicy(dlqPolicy) .subscribe();maxRedeliverCount表示最大重新投递次数这里设置 3 次。注意这个 3 是怎么计算的它会算上正常消费失败那一次吗严格来说这个值是指投递次数redelivery count。如果一条消息第一次消费失败触发 nack然后被重新投递了 3 次仍然失败那么就会进入死信队列。所以实际处理的尝试次数是 4 次。这里有个容易误解的概念设置enableRetry(true)之后如果又配了deadLetterPolicy重试和死信是同时生效的——先走重试重试次数达到maxRedeliverCount后进死信。如果只配置deadLetterPolicy而没开启 retry那么maxRedeliverCount也可以在非重试模式下生效只不过没有延迟排队的过程失败的投递会在negativeAckRedeliveryDelay的控制下重投直到达到上限。3.3 死信消息的消费与恢复死信主题里的消息不会凭空消失需要单独写一个消费者去接管。我的建议是死信消费逻辑必须做两件事第一是告警——补一个监控只要死信主题有新消息立刻报警第二是记录日志并触发人工介入流程。恢复死信消息时最常用的手段是读取死信消息的内容修复业务数据或补偿依赖方然后重新发送到原始主题里重新消费。很多人会直接把死信消息从 DLQ 重新发送到主队列但这有一个风险如果消息本身是“毒消息”比如 payload 格式不对重发多少次都会再次死信造成循环。所以恢复前一定要先看消息内容的失败原因最好是连重试几次都无法处理的原因。我自己的习惯是给死信消息打个标记在重新入队之前校验一下 payload 是否符合预期。如果内容本身有问题直接走人工修复流程而不是盲目重发。这比任何自动化策略都可靠。4. MessageId 深度解析为什么是 messageid|28077:20854:04.1 MessageId 由哪几部分组成有读者问过一个问题为什么 Pulsar 的 MessageId 看起来是messageid|28077:20854:0这种奇怪格式而不是像 Kafka 那样一个简单的offset123456。我来拆开讲。在 Pulsar 中MessageId是一个结构化的对象默认的toString()输出格式是ledgerId:entryId:partitionIndex对于非批消息如果有批次还可能包含batchIndex。拿28077:20854:0来说28077是 ledger ID20854是 entry ID最后一个0是 partitionIndex分区序号。那这背后对应什么存储结构呢Pulsar 的消息是存在 Apache BookKeeper 里的BookKeeper 按 ledger账本来组织数据。一个 ledger 可以理解成一段连续写入的日志文件集合消息写入时会被切分成 entry每条 entry 就是一个逻辑记录单元。所以28077指的是这条消息落在第 28077 个 ledger 上20854是它在 ledger 里的第 20854 条 entry。最后那个0是分区序号。Pulsar 的 topic 可以分区分区数决定了同一主题下消息分布到不同的分区。每个分区都有自己的 ledger 链所以消息 ID 最后一段通常就是 partition index。如果你消费的是单分区主题这里一般就是 0如果是多分区主题你可能会看到...:0、...:1这样的差异。4.2 为什么不是 Kafka 那样的单调递增 offset很多人从 Kafka 切到 Pulsar 后第一反应是Kafka 的 offset 就是一个分区内单调递增的 long 型数值查询和定位都方便。Pulsar 为什么搞得这么复杂这是存储引擎的差异造成的。Kafka 是本地分区日志模型每个分区的日志是连续的、可顺序读写的offset 天然就是一个递增偏移量。而 Pulsar 把存储层抽成了 BookKeeper消息写入时不保证顺序——不同生产者的消息会并发写入到不同的 ledger 里——所以需要ledgerId entryId这种二维坐标才能唯一定位一条消息。类比一下Kafka 的 offset 像图书馆里一本书的页码顺序是固定的Pulsar 的 MessageId 像快递的运单号ledgerId是哪个仓库发出的entryId是这个仓库里的第几单partitionIndex是走哪条运输线。这种设计让 Pulsar 在存储容量、水平扩展和故障恢复上有更大的弹性但也让消息定位变得不那么直观。在实际排障时这个差异很重要。你在日志里看到messageid|28077:20854:0就可以判断这是ledger 28077、entry 20854、partition 0的消息。如果想深入排查可以登录 BookKeeper 集群用bkctl或者直接查 ledger 元数据来定位这段数据。4.3 MessageId 与重试、死信的关系消息的 MessageId 在重试和死信流程里扮演了一个很关键的角色客户端向 Broker 发送nack时必须携带这条消息的 MessageIdBroker 才能知道要重新投递哪条消息。死信写入时同样基于 MessageId 的原始信息。这里还有一个容易踩坑的点批消息的 MessageId。Pulsar 支持生产者把多条消息打包成一个批次写入这样一个 entry 里可能包含多条消息。如果你处理的是批消息MessageId 的尾部可能不止三段而是类似ledger:entry:partition:batchIndex的四段。在重试和死信配置中你必须用MessageId 的 batchIndex来区分同一条 entry 里的不同子消息否则会误删或重复消费。我在一个生产案例里遇到过这种问题消费者 ack 了一条批消息用的是批消息的整体 ID结果整个批次的其余消息都被标记为已确认直接丢了。正确做法是一条一条地调用consumer.acknowledge(msg.getMessageId())或者用MessageIdAdv之类的工具判断是否带 batchIndex。这是 Pulsar 重试机制里的一个隐藏深坑面试和实战都会碰到。5. Pulsar 和 Kafka资料丰富度对比与学习路径5.1 生态和资料的体感差异有读者问“pulsar和kafka哪个资料丰富一些”我的体感是Kafka 的资料量级至少是 Pulsar 的五到十倍。很多东西的搜索结果Kafka 能翻到从入门到原理到源码解析的全栈内容而 Pulsar 往往会止步于官方文档和少部分博客。这不是说 Pulsar 不好而是发展阶段不一样。Kafka 从 2011 年在 LinkedIn 落地到现在十几年社区用户基数大市面上从《Kafka 权威指南》到各种付费课程都沉淀得很厚。Pulsar 是 2018 年前后才开始快速普及的很多资料还停留在基础用法层面深入分析重试机制、BookKeeper 存储细节的文章相对稀缺。如果你正在选型我建议不要纯看资料丰富度。Kafka 资料丰富是因为存量用户多、踩坑的人多Pulsar 资料少但架构特性更丰富比如存算分离、多租户、统一队列和流模型。如果你团队能沉下心啃官方文档和源码Pulsar 的坑并不比 Kafka 更多。关键是团队里有没有人愿意做那个“资料开拓者”。5.2 Pulsar 学习路径的实操建议既然资料稀缺那该怎么高效学习 Pulsar尤其是重试和死信这块我的路线是这样的第一步花三天时间把官方文档的 Consumer、Message、Retry 和 Dead Letter 四个章节啃一遍。官方文档是所有资料里最准确、更新最快的很多问题网上搜不到答案但文档里其实写了。我见过太多人遇到问题直接问群里而答案就安静地躺在官方文档某段小字里。第二步找一套真实的业务场景做边学边练。别只跑standalone模式从控制台发一条消息。你至少应该模拟一个消费者配置 nack、重试、死信然后在控制台观察消息的流转路径。没有生产环境的用 Docker 起一个单机 Pulsar 集群完全够用。第三步深入看源码。Pulsar 客户端源码不算太复杂NegativeAcknowledgement的处理逻辑、RetryTopic的实现都在org.apache.pulsar.client.impl包里。我花了几个晚上把NegativeAcksTracker和DeadLetterTopicPolicy相关代码过了一遍之后排障的底气完全不一样。5.3 从 Kafka 迁移过来需要注意的思维转变如果你是从 Kafka 转过来的有三个思维转变特别重要。第一是 offset 这个概念的弱化。Kafka 排障常说“consumer offset 落后了”“重置 offset”但 Pulsar 里没有全局 offset 的概念只有MessageId和Subscription的markDeletePosition。你不能像 Kafka 那样直接“跳到某条 offset”而是要通过seek(MessageId)来回放。例如consumer.seek(28077:20854:0)可以直接让消费者从这条消息重新开始消费这在修复重复消费或数据错误时非常常用。第二是 ack 的粒度与方式。Kafka 的 offset 提交可以做到批量、异步、自动Pulsar 也类似但 Pulsar 对单条消息的 ack 更精细化。特别是批消息场景必须注意 batchIndex。第三是消费订阅模型。Kafka 的 consumer group 在 Pulsar 里对应Subscription但 Pulsar 增加了Exclusive、Shared、Failover、Key_Shared四种模式。重试和死信在Shared、Key_Shared模式下行为差异很大这直接影响你配置重试策略时的判断。这块是文档之外最容易踩雷的地方。6. 生产配置实战重试 死信完整方案6.1 场景设定假设有一个订单服务消费者订阅orders主题处理订单时需要调用外部支付结果查询接口。接口偶尔超时且超时通常在 3 到 5 秒订单数据偶发格式异常。我们希望实现普通超时失败先重试 3 次每次间隔 10 秒重试后仍然失败进入死信队列格式异常的消息直接进死信队列不进行无效重试。这个场景非常典型既要给“暂时性失败”足够的恢复机会又要避免对“永久性失败”的无效折腾。下面用 Java 客户端来实现。6.2 完整配置与代码首先是消费者配置DeadLetterPolicy dlqPolicy DeadLetterPolicy.builder() .maxRedeliverCount(3) .deadLetterTopic(persistent://public/default/orders-dlq) .build(); ConsumerOrder consumer client.newConsumer(JSONSchema.of(Order.class)) .topic(persistent://public/default/orders) .subscriptionName(order-pay-query-sub) .subscriptionType(SubscriptionType.Shared) .enableRetry(true) .negativeAckRedeliveryDelay(10, TimeUnit.SECONDS) .deadLetterPolicy(dlqPolicy) .subscribe(); while (true) { MessageOrder msg consumer.receive(5, TimeUnit.SECONDS); if (msg null) continue; try { Order order msg.getValue(); // 处理订单查询支付结果 boolean ok processOrder(order); if (ok) { consumer.acknowledge(msg.getMessageId()); } else { consumer.negativeAcknowledge(msg.getMessageId()); } } catch (InvalidOrderException e) { // 消息本身有问题直接进死信不走重试 consumer.reconsumeLater(msg.getMessageId(), 1, TimeUnit.MILLISECONDS); // 但注意reconsumeLater 也是 reconsume // 如果 maxRedeliverCount 设置了达到后会进死信。 // 如果想让这种消息立即进死信需要在 catch 中调用 // consumer.reconsumeLater(msg, 0, TimeUnit.MILLISECONDS); // 并用 ControlMessage 或独立策略区分 } catch (Exception e) { consumer.negativeAcknowledge(msg.getMessageId()); } }这里有个细节InvalidOrderException和普通异常区分处理。但因为maxRedeliverCount的存在所有 nack 或 reconsumeLater 都会计入重试次数所以哪怕是格式异常的消息也会至少重试几次才进死信。如果你希望“格式异常立即进死信”就不能用同一个消费者或者得在消费者外部用一个校验层做前置过滤直接把不合格消息发送到死信主题。我在生产中的做法是消费者在反序列化阶段如果抛异常直接获取原始byte[]构造一条新消息投到 DLQ 主题然后 ack 掉当前消息相当于手动完成“立即死信”。为什么我这么在意立即死信因为格式异常消息的重试毫无价值它占用的每次重试都会给 CPU、网络和日志系统带来压力。你要做的不是让系统“多试几次”而是“准确地试几次”。6.3 验证与监控配置完成后验证三步走。第一步向主题发送一条会产生InvalidOrderException的消息确认它直接进入 DLQ第二步发送一条依赖接口超时的消息确认它先重试 3 次观察消费日志每隔 10 秒出现一次然后进入 DLQ第三步发送一条正常消息确认它首次处理成功并 ack。监控指标方面我特别推荐关注 Pulsar 控制台中的Subscription backlog、Unacked messages以及死信主题的消息数量。一旦死信数量开始上升大概率说明某个下游服务出了问题这个信号比业务日志更直接、更快。7. 常见问题与排查技巧实录7.1 问题速查表现象可能原因排查思路消费失败后消息不重试没开启enableRetry或 nack 没生效检查 ack 超时和 nack 参数确认调用了negativeAcknowledge()重试间隔不符合预期negativeAckRedeliveryDelay配置过短或过长检查两个配置nack 延迟和 ackTimeout死信消息迟迟不出现maxRedeliverCount过大或计数分散调小阈值检查共享订阅模式下的多实例计数消息进入死信但业务没告警没监控 DLQ 主题给 DLQ 主题增加消息积压监控批消息重复消费或丢失没有按 batchIndex 逐条处理使用MessageIdAdv解析批索引逐条 ack恢复消息后再次死信毒消息没有修复就重投校验 payload修复后再投递messageid 格式解析不明白不了解 ledger/entry/partition 模型按本文第 4 节的思路逐步定位这张表是我日常支持同事时最常用的排障路径。多数问题不是 Pulsar 本身出了 bug而是配置语义没吃透。7.2 独家避坑心得最后分享几个我反复踩过的坑希望能帮你少走弯路。第一个坑是关于 nack 和 ackTimeout 的联动。我在一个服务里只设置了negativeAckRedeliveryDelay(5s)但没有设置ackTimeout。结果某次消费者线程卡死消息既不 ack 也不 nack就一直卡在那。后来把ackTimeout设置为 15 秒并配置了ackTimeoutRedeliveryBackoff才避免这类“卡死不自动恢复”的问题。第二个坑是共享订阅模式下的重试计数分散。我在双实例部署下明明设置了maxRedeliverCount3但消息重试了 6 次才进死信。原因就是两个实例各自计数。解决方案把重试次数的判定做成“基于延迟与业务表记录”的幂等策略或者在业务侧记录这条消息的处理失败次数达到阈值后主动 ack 并投递到 DLQ。第三个坑是关于 Key_Shared 订阅下的重试主题行为。Key_Shared 模式下重试消息会按 key 保证顺序这是一把双刃剑一个 key 的消息失败了能保证同 key 后续消息不抢跑但如果这个 key 一直失败会阻塞同 key 的所有后续消息。我在设计通知类业务时特意用了SubscriptionType.Shared来避免这种头阻塞问题。第四个坑是死信主题的自动创建。早期版本里如果你没有手动创建 DLQ 主题且客户端没有 auto-create 权限死信消息投递会失败。现在新版本默认允许客户端自动创建主题但生产环境出于安全限制通常会关掉allowAutoTopicCreation。所以配置完死信后记得提前确认orders-dlq主题存在否则进死信时消息会直接报错反而变成“死信也死了”。这个细节很多文档都没写但线上就是这么残酷。个人经验来说Pulsar 的消息重试和死信机制在功能完备性上是我用过几款消息队列里最灵活的。但灵活性也意味着复杂度。你不需要一开始就把所有参数调到完美但一定要先保证“失败不丢消息、重试有上限、死信有监控”这三条底线。踩过这几个坑之后再回头看官方文档里那些参数你会觉得每一行都特别顺眼。
分享:

看完干货,该让你的企业上线了

免费需求沟通 · 48 小时内出具建站方案 · 河南本地可上门