ARTICLE DETAIL

资讯详情

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

RabbitMQ进阶实战:死信队列、延迟队列与消息防丢失全解析

RabbitMQ进阶实战:死信队列、延迟队列与消息防丢失全解析 先说我为什么想写这篇东西。前两天线上一个订单超时未支付自动关闭的功能突然不执行了查了半天发现不是代码逻辑的问题而是原来的延迟消息方案在服务重启后丢了一批消息后续补偿又没做好导致一批订单一直挂着没人管。这种问题只要用过RabbitMQ做业务的人早晚都会遇到——死信队列怎么用、延迟队列怎么做、消息怎么才能不丢这三件事如果只是在文档层面知道个大概一上生产就会踩坑。所以这篇我打算把RabbitMQ进阶里最关键的三块放在一起讲清楚死信队列的触发机制和实际业务映射、延迟队列的两种实现路径到底怎么选、防丢失机制从生产者到消费者的每一环应该怎么配置。内容偏实战适合已经会用RabbitMQ收发消息、正在琢磨怎么把它用到真实业务场景中的开发者。1. 死信的本质一条消息是怎么被“判死刑”的很多初学者把死信队列当成一个独立的组件来学其实理解反了。死信队列在RabbitMQ里不是什么特殊队列它就是一个普普通通的队列只是这个队列专门用来接收那些在别的队列里“待不下去”的消息。所以问题的关键不在于死信队列本身而在于一条消息在什么情况下会被判定为死信以及被判定之后它会被路由到哪里。1.1 死信的三个来源TTL过期、队列溢出、消费者拒绝RabbitMQ里一条消息成为死信官方定义有且只有三种情况。第一种是消息过期。你可以给队列设置一个x-message-ttl参数也可以给单条消息设置过期时间消息在队列里待了超过这个时间还没被消费就会被标记为过期。这条消息一旦过期如果队列没有配置死信交换器它会被直接丢弃听起来像没什么问题但在很多业务里一个突然消失的消息可能意味着一次超时订单没被关闭、一个优惠券没被作废。所以“过期消息不丢弃、而是转到死信队列”这个能力是延迟队列实现的基石。第二种是队列达到最大长度。队列可以设置x-max-length或者x-max-length-bytes当队列里的消息数量超过这个限制时新来的消息会被拒绝进入或者按策略从队头丢弃老消息。默认情况下被丢弃的消息就没了但是如果配置了死信交换器这些溢出消息也会变成死信送到死信队列去。第三种情况是消费者主动拒绝。消费者拿到消息后如果业务处理失败可以调用basicNack或者basicReject并且在参数里指定requeuefalse消息就不会再回到原队列而是进入死信队列。如果你指定requeuetrue消息会被重新放回原队列继续投递这种情况下不会触发死信。这三条来源在业务里各有各的映射TTL过期最常见用来做订单超时、延迟回调队列溢出的场景没那么频繁但流量突然暴涨、消费者大面积卡死的时候溢出策略能挡住一部分冲击配合死信队列可以做到“丢消息也要丢得明明白白”消费者拒绝加requeuefalse则适合做业务上明确失败、重试也没有意义的消息比如库存不足、参数校验不通过这种消息与其反复重投不如直接送去死信队列由人工或者兜底程序去处理。1.2 死信去向x-dead-letter-exchange 与 x-dead-letter-routing-key 的配合逻辑光知道消息会变成死信还不够你还得告诉RabbitMQ“死信往哪送”。这就要在声明队列的时候加上两个参数Bean public Queue orderDelayQueue() { return QueueBuilder.durable(order.delay.queue) .withArgument(x-message-ttl, 30000) .withArgument(x-dead-letter-exchange, order.exchange) .withArgument(x-dead-letter-routing-key, order.expire) .build(); }x-dead-letter-exchange指定死信要发送到哪个交换器x-dead-letter-routing-key指定发送时用的路由键。这里有一个细节容易忽略如果x-dead-letter-routing-key没有显式配置RabbitMQ会默认使用这条消息进入当前队列之前的原始路由键。也就是说如果发送消息时用的路由键恰好能被死信交换器正确匹配你甚至不用配死信路由键。这种转发逻辑在我看来有点像快递派送失败后的“改址投递”。原地址原队列派送不了快递员RabbitMQ会根据你填的单子把这件包裹重新分配到另一个网点死信队列。要是你没填改址信息包裹就只能退回来或者丢弃。还有一点得提醒没有配置死信交换器的队列在消息过期、溢出或拒绝时消息会被静默丢弃。这在开发环境问题不大但在生产环境等于给自己埋雷。我见过不止一次线上消息突然少了几条查了半天发现是消费异常时requeue设了false但队列根本没配死信交换器消息直接被丢了。2. 延迟队列的两种实现TTL死信叠加与官方延迟插件先纠正一个常见误解RabbitMQ本身并没有一个叫“延迟队列”的队列类型。你搜遍官方文档也找不到一个叫delayed-queue的属性。所谓延迟队列是通过“TTL过期 死信转发”的组合能力变相实现的或者借助官方提供的延迟消息插件来实现。2.1 方案一TTL死信零插件实现延迟调度这个方案的原理很简单生产者不直接把消息发到业务消费队列而是先发到一个设置了TTL的队列这条消息在里面“睡”指定时间睡醒之后过期变成死信被RabbitMQ自动转发到真正处理业务的死信队列。整个链路是生产者 → 延迟队列带TTL → 死信交换器 → 业务队列 → 消费者。以订单超时关闭为例声明两个队列一个延迟队列用于缓冲一个业务队列用于真正处理Bean public Queue delayQueue() { return QueueBuilder.durable(order.delay.queue) .withArgument(x-message-ttl, 30000) .withArgument(x-dead-letter-exchange, order.exchange) .withArgument(x-dead-letter-routing-key, order.close) .build(); } Bean public Queue closeOrderQueue() { return QueueBuilder.durable(order.close.queue) .build(); } Bean public DirectExchange orderExchange() { return new DirectExchange(order.exchange, true, false); } Bean public Binding closeOrderBinding() { return BindingBuilder.bind(closeOrderQueue()) .to(orderExchange()) .with(order.close); }生产者正常发送到order.delay.queue30秒后消息过期自动流转到order.close.queue消费者监听到之后执行关单逻辑。这套方案的好处是零插件、零额外依赖用的是RabbitMQ原生能力任何环境都能跑。但这个方案有一个天然缺陷必须提前知道TTL不保证时效精确。RabbitMQ对过期消息的检查逻辑是惰性的只有消息到达队列头部时才会检查它是否过期。如果一个延迟队列里有100条消息队首那条TTL是30秒后面的某条TTL是5秒那么后面的消息必须等队首消息先被处理完或者先过期才会轮到自己。也就是说短TTL消息可能被前面的长TTL消息阻塞实际延迟时间远大于预期。这个问题在做秒级精确延迟业务时基本不可接受。另外使用TTL死信方案时业务队列其实收到的是“延迟队列里原来的那条消息”。如果你需要感知消息在延迟队列里待了多久原始的消息收发明细和存活时间会以死信头的方式携带但读取和处理都更绕。对大多数业务来说这一层不用关心但排查问题的时候知道这个消息字段能从死信头里取会省不少事。2.2 方案二延迟插件秒级精确触发的适用范围为了解决TTL方案的时效精度问题RabbitMQ官方提供了rabbitmq_delayed_message_exchange插件。安装插件后交换器多了一种类型叫x-delayed-message消息发到这种交换器后不会立刻路由到队列而是由插件延迟器托管等到延迟时间到达后再完成真正的路由。启用插件之后声明一个延迟交换器Bean public CustomExchange delayedExchange() { MapString, Object args new HashMap(); args.put(x-delayed-type, direct); return new CustomExchange(order.delayed.exchange, x-delayed-message, true, false, args); } RabbitListener(bindings QueueBinding( value Queue(value order.queue, durable true), exchange Exchange(value order.delayed.exchange, type x-delayed-message, delayed true), key order.close )) public void handleCloseOrder(OrderMessage message) { // 业务处理 }发送消息时在消息属性里带一个x-delay的header单位是毫秒插件会根据这个值决定延迟多久再投递。从实测效果看这个方案在秒级以下精度也能保持较好的可靠性而且不会出现TTL方案里队头阻塞的问题每条消息的延迟时间是独立计算的。所以在订单超时、定时红包提醒、延迟通知这类对时间精度有要求的场景我建议优先用插件方案。当然插件方案也有成本一是需要额外安装插件并保证插件版本和RabbitMQ主版本兼容升级RabbitMQ时容易踩坑二是它本质上改变了消息的投递路径调试的时候需要多一层认知负担。下面这个表是两种方案的对比对比项TTL死信延迟插件额外依赖无需要安装插件版本需匹配延迟精度受队头阻塞影响不精确每条消息独立计时秒级精度适用场景分钟级以上的业务延迟秒级延迟、时间敏感业务运维成本低原生能力中需关注插件升级消息可追溯性消息会经过死信流转头信息变化交换器内部托管对消费者无感如果业务场景只是“30分钟后关单”“24小时后提醒”这种TTL死信完全够用没必要为了这些场景引入插件。但如果你的延迟业务时间跨度小、精度要求高或者存在大批量不同延迟时间的消息混杂在同一个队列里别犹豫直接上插件方案。3. 防丢失机制三段链路每一跳都要显式守护RabbitMQ消息丢失的问题十个生产事故里至少占三个。很多时候不是因为RabbitMQ本身设计有问题而是使用者默认“消息发出去就等于消息到了”。从生产者到消费者消息其实要经过三段链路生产者 → 交换器 → 队列 → 消费者任何一跳出了问题消息都会悄无声息地丢。我强烈建议做消息防丢失时别总想着有什么“一键开启”的开关一段一段往下查每一段都确认到位才能真正做到不丢。3.1 第一跳生产者confirm机制的三个确认级别生产者发送消息后Broker收没收到RabbitMQ默认是不告诉你的。好在RabbitMQ提供了publisher confirm机制开了之后Broker每收到一条消息都会异步回调一个确认结果。Spring Boot里配置非常直接spring: rabbitmq: publisher-confirm-type: correlated publisher-returns: true然后给RabbitTemplate设置回调rabbitTemplate.setConfirmCallback((correlationData, ack, cause) - { if (ack) { log.info(消息确认成功: {}, correlationData.getId()); } else { log.error(消息确认失败: {}, cause: {}, correlationData.getId(), cause); } }); rabbitTemplate.setReturnsCallback(returned - { log.error(消息未路由到队列: {}, exchange: {}, routingKey: {}, returned.getMessage(), returned.getExchange(), returned.getRoutingKey()); });这里有两层回调要分清楚ConfirmCallback管的是“Broker收到消息没有”ReturnsCallback管的是“消息到底有没有路由到队列里”。比如消息发到了交换器但交换器根据路由键找不到任何队列这种场景下confirm可能返回了成功但消息实际上没有被任何队列接收。只有两个回调都正常消息才算真正落到了队列里。生产环境做发送端防丢失我的习惯是确认失败的消息不直接丢弃而是存入一张本地消息表由定时任务定期重投。配套的逻辑类似一个两阶段提交记录消息状态为发送中、发送成功或发送失败定期把发送失败且超过重试次数的消息捞出来人工处理永远不会因为Broker抖动而把消息搞丢。3.2 第二跳交换器、队列、消息三者的持久化组合消息到达队列后如果此时RabbitMQ服务器宕机了消息还在不在这取决于三个东西是否都设置了持久化缺一不可。第一个是交换器持久化。声明交换器时有durable参数设置为true才能保证交换器在Broker重启后仍然存在否则交换器会消失之后无法继续接收消息。第二个是队列持久化QueueBuilder.durable(xxx)否则队列重启后也没了消息自然跟着消失。第三个是消息持久化设置MessageDeliveryMode.PERSISTENT在Spring Boot里默认就是持久化模式但如果你用原生客户端或者手动构造Message一定记得显式设置。只有这三者同时满足一条消息才算真正落到了持久化存储上。但即便这样依然存在一个极小概率的丢消息窗口消息已经被标注为持久化但还没来得及刷到磁盘时Broker断电了。这种情况严格来说无法完全避免RabbitMQ官方对此也明确过如果业务对“绝对不丢”有极高要求建议使用Quorum Queue仲裁队列它通过多副本复制实现高可用比经典队列的持久化更稳。不过Quorum Queue对吞吐量有一定损耗加入之前要做压测确认业务能够接受。我自己的建议是绝大多数业务配上confirm 持久化 手动ack已经足够仲裁队列除非是金融级、对账级的消息链路否则没必要轻易引入。3.3 第三跳消费端手动ack与幂等兜底消息从队列投递给消费者之后如果消费者处理到一半挂掉了这条消息呢这完全取决于消费者用的是自动ack还是手动ack。自动ack模式下Broker把消息发给消费者就立即认为“消息已被消费”哪怕业务逻辑还没执行这条消息也会从队列里移除。如果消费者处理业务抛异常消息已经回不来了。所以我每次都要强调涉及关键业务的消息一律关闭自动ack改成手动ack。Spring Boot里手动ack的配置稍微绕一点RabbitListener(queues order.queue, ackMode MANUAL) public void handleMessage(OrderMessage message, Channel channel, Header(AmqpHeaders.DELIVERY_TAG) long deliveryTag) { try { orderService.handle(message); channel.basicAck(deliveryTag, false); } catch (Exception e) { channel.basicNack(deliveryTag, false, false); } }basicAck确认消息处理成功basicNack的第三个参数表示是否重新放回队列。设置false就是“处理失败且不重回队列”这条消息会变成死信方便后续集中排查如果设置true消息会被重新投递配合重试机制可以在短时间内重新处理但如果业务本身有问题它就会一直在消费者和队列之间来回横跳形成无效重试循环。还有一层几乎人人都知道、但经常做不好的防护——幂等性。前面已经说了RabbitMQ无法保证消息100%不重复投递比如消费者处理成功但ack网络丢了Broker会重新投递这条消息所以消费端必须处理重复消息。最简单的做法是给每条消息一个全局唯一ID消费者拿到之后先查Redis记录如果已经处理过就直接ack不再执行重复逻辑。数据库层面也可以加唯一索引哪个顺手用哪个但这条必须有不然等到线上真的出现重复数据再来补就已经晚了。4. 排障手记三个真实踩坑案例理论说再多不如把线上的坑一个个摊开来看。这里分享三个我实际踩过、也帮别人排查过的问题每个都跟前面的内容直接相关。4.1 死信循环路由键配置失误造成消息“鬼打墙”有一次收到监控告警某个死信队列的消息量在短时间内暴涨消息堆积到几十万。第一反应是业务消费失败看了一下死信队列里的消息内容发现都是同一批订单超时消息而且这些消息被反复转发A队列过期 → 死信转到B队列 → B队列又配置了死信转发到A队列 → 又过期 → 又转回。这就是典型的死信循环。出现这种问题的原因通常是声明队列时两个队列各自配置了x-dead-letter-exchange而且两个交换器的绑定关系恰好形成了闭环。排查链路比较简单但我发现很多人在新建队列时根本没意识到“死信队列本身也可以配置死信转出”等到出了问题才回头检查。处理这种问题套路是线上排除掉积压消息之后把死信链路的配置梳理成一张拓扑图确保它是单向的、有终点的。死信队列作为兜底不要再给它配置下跌交换器。必要的时候用x-max-length限制死信队列的最大消息数防止积压风险滚雪球。4.2 手动ack漏写Redis里堆了两万条未确认消息还有一个更隐蔽的坑。之前有位同事改代码在RabbitListener方法里加了大量的业务校验逻辑中间写了好几个分支提前return但那些提前return的分支都没有执行basicAck也没有basicNack。消息被消费者拿走了但一直没有确认Broker认为消息还在消费者手里不会重投、不会超时删除就一直挂在Unacked状态。服务运行了几天之后Unacked消息堆到了两万条队列的Ready消息数却是0看起来一切正常实际业务早已停滞。排查这问题时我的建议是凡是开启手动ack的监听方法一定把basicAck和basicNack放在统一的入口而不是每个分支都写一遍。最稳妥的办法是使用Spring的RabbitListenerErrorHandler配合AOP统一处理异常或者干脆在方法级别用try-finally包裹确保任何路径下都会对消息做出ack或nack响应。另外监控上要盯住Unacked指标。这个指标异常太有迷惑性了——队列表面没有积压消费者也“看似”正常实际上消息全卡在确认环节上。4.3 TTL优先级队列级TTL和消息级TTL同时存在时的坑使用TTL死信方案做延迟队列时可能同时遇到队列级TTL声明队列时的x-message-ttl和消息级TTL发送单条消息时设置的expiration属性。一旦两者都设置了RabbitMQ的规则是取两者中较小的那个值作为过期时间。这个逻辑本身不难理解但实际问题往往出在消息级TTL的格式上。expiration属性接收的是字符串比如30000代表30秒。如果你不小心传了30000以外的格式比如30s这种人类可读格式RabbitMQ不会报错但这个字符串在broker端解析时会出问题导致这条消息的过期时间“看起来没有设置”实际上走了队列级TTL。这类问题时有时无排查起来很费劲。所以我的建议是延迟队列业务中统一使用队列级TTL消息级TTL除非有充分的理由否则别混用避免优先级和格式带来的双重困扰。5. 生产落地清单参数、监控与上线前自检到了最后落地这一步很多细节必须文档化。项目换人、多团队协作的时候如果没有一份清晰的参数和监控清单早晚会有人因为配置错误引发线上问题。这里分享一份我自己的基线配置可以直接抄。5.1 一套可直接照抄的声明参数清单先说延迟队列场景的完整参数参数推荐值说明x-message-ttl按业务设定单位ms延迟时间仅TTL方案需要x-dead-letter-exchange业务实际交换器名必配否则过期消息直接丢弃x-dead-letter-routing-key业务队列绑定的路由键不配默认沿用原路由键x-max-length按业务积压容忍度定防止队列无限堆积压垮内存durabletrue队列持久化必须开启auto-deletefalse集群中队列自动删除会造成诡异问题消息投递模式PERSISTENTSpring默认已设置良好的机制生产者的confirm配置spring: rabbitmq: publisher-confirm-type: correlated publisher-returns: true消费者的手动ack配置spring: rabbitmq: listener: simple: acknowledge-mode: manual retry: enabled: true max-attempts: 3 initial-interval: 1000上面这个重试配置是消费者内部的重试重试超过最大次数后消息会被丢弃或进入死信取决于异常类型。注意它和basicNack的requeue重投是两个不同层面的重试别搞混。5.2 监控指标与告警阈值建议光把参数配好还不够没有监控等于睁着眼睛开车。结合我自己的经验这几个指标的告警阈值直接按下面的方式设定能帮你大多数情况下提前发现险情指标含义建议告警阈值Ready待消费消息数持续5分钟超过业务正常水位如1万Unacked已投递未确认消息数持续5分钟超过100基本等于消费卡死死信队列incoming速率每秒进入死信队列的消息数超过正常速率2倍队列总消息深度Ready Unacked超过设定上限即告警连接数当前TCP连接数异常陡增或陡降都要查一个重要心得监控告警不能只看队列积压数字一定要看“死信队列的引入速率”。很多时候业务故障已经发生了但Ready消息数还在正常范围因为问题消息全被打进了死信队列。不盯死信速率等于把最大的报警信号漏掉了。上线前自检还有一个容易被忽略的点——消费端的幂等性。建议在所有消息消费入口做一个统一的幂等拦截不管RabbitMQ怎么重投、重试业务侧永远只处理一次。这个底兜做好了前面说的任何防丢失机制就算配置失误至少不会产生脏数据和重复数据。最后再说一点个人习惯。每次我在配置死信队列时都会专门用一个队列名称后缀_dlq并且不允许在死信队列上再配置死信交换器从根上断绝死信循环的可能。这几条经验看着零散但实际使用中能帮你少走非常多弯路。
返回列表