做大数据链路的朋友应该都有过这种经历:凌晨三点,数据同步任务突然断流,上游日志照常生产,下游数仓却几个小时没有新分区。查来查去,问题出在消息队列——某个消费端抛了个空指针,消息没确认,队列里的消息越积越多,最终引爆了堆积告警。这种时候,RabbitMQ的消息重试机制就是你最需要掌握的救命技能。
这篇文章我会围绕大数据场景下RabbitMQ消息重试机制的实现思路,从方案选型、核心原理、代码实践到问题排查,完整过一遍。内容包括自动重试与手动确认的取舍、死信队列和延迟重试的搭建、重试参数的计算逻辑,以及我在高吞吐数据管道里实测下来的踩坑经验。适合正在做数据同步、订单处理、日志清洗这类实时链路的朋友,也适合想系统搞懂RabbitMQ重试机制的新手。
1. 为什么大数据场景离不开消息重试机制
1.1 没有重试机制,大数据管道会怎样
先把问题说透。大数据场景下,消息生产端通常是日志采集组件、业务库binlog同步工具或者实时计算引擎,消费端往往是清洗服务、转换服务、落库服务。这条链路上任何一个下游抖动,都会直接导致消息处理失败。最典型的就是下游HBase集群在做Region分裂,或者ES集群正在段合并,写入延迟飙到几秒甚至超时,消费端一次性拉取几百条消息批量写入,失败是常态而不是意外。
如果不做重试,会出现两种极端情况。第一种是消息直接丢弃,数据永久丢失,数仓第二天对不上数,补数据能补到怀疑人生。第二种是无限requeue,消息被退回队列头部,同一个消费端再次拉取再次失败,形成死循环,整条链路被一条脏数据卡死,后面所有消息全部积压。我见过一个真实案例,一条JSON里某个字段偶发为空导致解析异常,消费端每拉取一次就抛一次异常,队列堆积从几百涨到几百万只用了不到两个小时。
更隐蔽的问题是重试失败后的消息归宿。如果不设计死信队列,重试多次仍然失败的消息最终会被丢弃或者无限重投,这两种结果在数据一致性上都是灾难。大数据链路对数据完整性要求极高,每条消息都可能代表一笔订单、一次用户行为、一条设备日志,丢了再想补回来,成本远高于在消息队列层面把重试机制做完善。
1.2 方案选型思考:为什么是RabbitMQ
谈到消息队列选型,很多人第一反应是Kafka,因为大数据场景Kafka的出场率确实高。但RabbitMQ在重试机制这个维度上,反而有自己的独特优势。RabbitMQ的消息确认机制做得非常精细,支持自动ack、手动ack、nack、reject,配合死信交换机(DLX)、TTL、优先级队列这些特性,可以实现非常灵活的重试策略。Kafka的offset提交机制决定了它在重试场景下相对粗糙,要么暂停消费要么跳过,做不到RabbitMQ这种“消息级别”的精细控制。
RabbitMQ的重试机制核心由三部分构成:消费端重试配置、死信队列、延迟重试策略。这三者互相配合,能覆盖大多数失败场景。消费端重试解决瞬时故障(下游抖动、网络闪断),死信队列解决多次重试仍失败的消息的最终归宿问题,延迟重试解决“过一会儿再试可能就好了”的场景(如下游服务重启、依赖资源临时不可用)。
我当时选择RabbitMQ还有一个重要考量:团队里Java技术栈为主,Spring Boot对RabbitMQ的封装非常成熟,spring-boot-starter-amqp几乎零成本上手,配置项丰富,团队不需要额外学习成本。做技术选型不能只看性能指标,团队能撑起来、运维能接得住,才是真正的务实选择。
2. 消息重试的三种核心实现路径
2.1 自动确认与手动确认的本质区别
RabbitMQ的消费确认模式直接决定了重试机制的写法。自动确认模式下,消息一旦投递给消费者就被认为处理成功,Broker立刻删除消息。这种模式适合处理逻辑极其简单、失败概率极低的场景,比如纯粹的消息转发。但大数据场景下我不建议用自动确认,因为消费端处理批量数据、调用下游接口、写外部存储都可能失败,一旦自动确认,失败消息无法重新投递,相当于没有重试可言。
手动确认模式则是把主动权握在消费端手里。处理成功调用basicAck主动告知Broker删除消息;处理失败可以调用basicNack或者basicReject,告诉Broker消息处理失败,需要进入重试流程。手动确认是重试机制的地基,没有这个基础,后面所有策略都是空中楼阁。
有个细节很多人会忽略:basicNack和basicReject都有一个requeue参数。requeue设为true,消息会被放回原队列,可以再次被消费,但这是一种“无脑重试”,重试次数和间隔都不受控制。requeue设为false,消息会被路由到死信交换机,进入死信队列,等待后续处理。合理的设计通常是把两者结合起来:前几次重试走requeue快速重试,超过阈值后requeue=false进死信队列做兜底。
2.2 死信队列:重试失败的最终归宿
死信队列的机制值得多说几句。RabbitMQ在创建队列时可以声明x-dead-letter-exchange和x-dead-letter-routing-key两个参数。当消息满足死信条件时,Broker会把它重新发布到指定的死信交换机,再由死信交换机根据路由键投递到对应的死信队列。
触发死信的条件有三种:消息被消费端拒绝(basicReject或basicNack且requeue=false)、消息TTL过期、队列达到最大长度。在大数据场景下,最常见的是第一种。消费端在重试达到上限后,显式拒绝消息并指定requeue=false,消息就自动进入死信队列。
死信队列的价值在于把“处理不了的消息”和“正常业务消息”隔离开。死信队列可以单独建消费者,做人工介入处理、日志记录、异常监控,或者等下游恢复后重新投递。我在实际项目中会为死信队列配置独立的告警,死信队列一旦有消息进入就立刻报警,因为这通常意味着有系统性问题,而不是偶发故障。
2.3 消费端重试的关键参数与配置
Spring Boot环境下,消费端重试的配置非常简洁。在application.yml里可以这样配置:
spring: rabbitmq: listener: simple: acknowledge-mode: manual retry: enabled: true max-attempts: 3 initial-interval: 1000 multiplier: 2 max-interval: 10000这里有几个参数需要解释一下。acknowledge-mode必须设为manual,这是手动确认的开关。retry.enabled开启后,Spring会在监听器内部拦截异常并触发重试。max-attempts是最大尝试次数,注意这个值包含了第一次投递,所以配置3意味着第一次失败后再重试2次。initial-interval是首次重试的等待时间,multiplier是退避倍数,每次重试的间隔会按倍数递增,避免集中重试造成下游压力。
这套配置背后其实是Spring Retry框架在工作,它会对Listener的执行过程做AOP拦截,捕获异常后按退避策略重新执行。但这种重试有个限制:消息还处于未确认状态,没有真正退回Broker,重试期间消息一直在消费者内存里。如果消费者进程重启,这个消息就丢了。所以大数据量场景下,如果追求更高的可靠性,建议不要依赖框架重试,而是采用手动nack+requeue的方式,让Broker来管理重试。
3. 大数据场景下的完整重试架构设计
3.1 整体链路设计:从生产到死信的全流程
先说结论,我在大数据管道里用到的重试架构由三条队列组成:主队列、重试队列、死信队列。主队列承接上游消息,正常消费;消费失败的消息通过延迟插件或者TTL机制进入重试队列,等待一段时间后再投递回主队列;超过最大重试次数的消息进入死信队列,由专门的兜底消费者处理。
具体流程如下:上游服务把消息发布到主队列,主队列的消费者负责业务处理。处理成功后basicAck确认。处理失败时,第一次和第二次失败通过basicNack+requeue=false把消息投递到重试队列,重试队列设置了TTL(比如30秒),消息过期后自动重回主队列,消费者再次尝试处理。累计消费次数用消息头或者Redis计数,超过3次仍然失败,就把消息投递到死信队列。
这种设计的优势很明显:重试过程不占用消费者线程,不会因为某条消息卡住而阻塞后面消息的消费;重试间隔通过TTL控制,灵活且清晰;重试消费和生产解耦,主队列不会因为重试流量而堆积。整个链路把“处理”和“重试”两种行为从空间和时间上隔离开,适合大数据场景下的高吞吐需求。
3.2 重试次数与退避策略的计算逻辑
重试参数的设置我一直强调要算,不要拍脑袋。以订单数据同步场景为例,下游MySQL偶尔主从切换导致写入失败,通常30秒内能恢复;但如果下游磁盘满了或者连接池耗尽,可能需要几分钟甚至更久。重试间隔设计就要覆盖这两个场景。
我的经验公式是这样的:总重试时间窗口 = initialInterval * (multiplier^maxAttempts - 1) / (multiplier - 1)。比如initial-interval=1秒,multiplier=2,max-attempts=5,那么总时间大约是1+2+4+8+16=31秒。这适合快速自愈的故障场景。如果希望覆盖更长的故障窗口,可以把initial-interval调大或者max-attempts增多,比如 initial=5秒,multiplier=3,max-attempts=4,总时间就是5+15+45+135=200秒,大约3分多钟。
重试次数也不是越多越好。每次重试都会占用系统资源,消息在队列和消费者之间反复横跳也会增加网络开销。过高的重试次数意味着消息长时间滞留,造成下游数据延迟。做数据管道要明确一个原则:重试只是给下游恢复争取时间,不是解决根本问题的手段。超过3次仍然失败,说明大概率是程序Bug或者配置问题,继续重试只是浪费资源,不如进死信队列让开发排查。
3.3 幂等消费与重试的协同关系
重试机制必然带来一个问题:重复消费。因为消息可能在被确认之前重新投递,消费端必须做好幂等处理,否则一条订单会被重复写入两次,产生脏数据。
幂等方案常用的有三种:数据库唯一键约束、Redis分布式锁、业务状态机校验。我在订单同步场景用的是唯一键方案,在目标表设置业务订单号的唯一索引,重复插入时捕获DuplicateKeyException,当作成功处理即可。这种方案实现简单、性能好,适合高并发写入。
Redis方案适合判断“是否已经处理过”这类场景,处理前先SETNX一个带业务ID的key,处理成功后删除。Redis方案要注意TTL设置,TTL太短可能在重试时锁已经过期,起不到幂等作用;TTL太长又占用内存。我在实践中倾向把幂等和重试分开看待:重试解决“没成功”的问题,幂等解决“成功但重复”的问题,两者配合才能真正保证消息只被处理一次。
4. 高吞吐下重试机制的代码实现与踩坑实录
4.1 核心代码结构:手动确认+死信路由
下面是我在项目里落地的一套代码结构,核心点在于把确认逻辑和重试逻辑显式地写在业务代码里,不依赖框架隐藏的行为。
首先是队列声明,通过Java配置创建主队列、重试队列、死信队列以及对应的交换机:
@Configuration public class RabbitMQConfig { public static final String MAIN_QUEUE = "data.main.queue"; public static final String RETRY_QUEUE = "data.retry.queue"; public static final String DEAD_QUEUE = "data.dead.queue"; public static final String MAIN_EXCHANGE = "data.main.exchange"; public static final String DEAD_EXCHANGE = "data.dead.exchange"; @Bean public Queue mainQueue() { return QueueBuilder.durable(MAIN_QUEUE) .deadLetterExchange(DEAD_EXCHANGE) .deadLetterRoutingKey(DEAD_QUEUE) .build(); } @Bean public Queue retryQueue() { return QueueBuilder.durable(RETRY_QUEUE) .deadLetterExchange(MAIN_EXCHANGE) .deadLetterRoutingKey(MAIN_QUEUE) .ttl(30000) .build(); } @Bean public Queue deadQueue() { return QueueBuilder.durable(DEAD_QUEUE).build(); } @Bean public DirectExchange mainExchange() { return new DirectExchange(MAIN_EXCHANGE); } @Bean public DirectExchange deadExchange() { return new DirectExchange(DEAD_EXCHANGE); } @Bean public Binding mainBinding() { return BindingBuilder.bind(mainQueue()).to(mainExchange()).with(MAIN_QUEUE); } @Bean public Binding retryBinding() { return BindingBuilder.bind(retryQueue()).to(deadExchange()).with(RETRY_QUEUE); } @Bean public Binding deadBinding() { return BindingBuilder.bind(deadQueue()).to(deadExchange()).with(DEAD_QUEUE); } }这个配置的逻辑是:主队列绑定死信交换机,消息失败且requeue=false时进入死信交换机,死信交换机根据路由键把消息投递到重试队列或者死信队列。重试队列设置了TTL为30秒,消息在重试队列里待满30秒后,自动路由回主交换机,再次进入主队列被消费。注意这里的主队列需要绑定主交换机,而重试队列的死信交换机配置为主交换机。
死信路由键的设计需要一致:主队列声明deadLetterRoutingKey为DEAD_QUEUE,这样失败的消息会进入死信队列;而如果我希望失败后先进重试队列而不是直接进死信队列,严格来说是需要两个不同的失败策略来区分。这里我用了更灵活的方式:消费端在nack时通过basicPublish把消息重新发送到重试队列,而不是依赖死信路由,避免不同路由策略互相打架。
4.2 消费端手动确认的完整写法
消费端的核心逻辑包括处理、确认、重投、进死信四步。我的实现思路是记录消息的重试次数,用消息头retryCount来标识,每次重试加一,超过阈值就拒绝进入死信:
@Component @Slf4j public class DataConsumer { private static final int MAX_RETRY_TIMES = 3; @Autowired private RabbitTemplate rabbitTemplate; @RabbitListener(queues = RabbitMQConfig.MAIN_QUEUE) public void onMessage(Message message, Channel channel) throws Exception { long deliveryTag = message.getMessageProperties().getDeliveryTag(); String msgBody = new String(message.getBody()); int retryCount = getRetryCount(message); try { processData(msgBody); channel.basicAck(deliveryTag, false); } catch (Exception e) { log.error("消息处理失败,业务数据={},重试次数={}", msgBody, retryCount, e); if (retryCount >= MAX_RETRY_TIMES) { // 超过重试上限,进入死信队列 channel.basicNack(deliveryTag, false, false); } else { // 增加重试次数,发送到重试队列 MessageProperties props = MessagePropertiesBuilder.newInstance() .setHeader("retryCount", retryCount + 1) .setExpiration("30000") .build(); Message retryMsg = MessageBuilder.createMessage(message.getBody(), props); rabbitTemplate.send(RabbitMQConfig.DEAD_EXCHANGE, RabbitMQConfig.RETRY_QUEUE, retryMsg); channel.basicAck(deliveryTag, false); } } } private int getRetryCount(Message message) { Object count = message.getMessageProperties().getHeader("retryCount"); return count == null ? 0 : (int) count; } private void processData(String data) { // 模拟大数据处理流程:解析、转换、写入下游 } }这里有几个处理细节值得说明。一个要点是,当业务逻辑失败需要重试时,我先通过rabbitTemplate发送一条带retryCount头的新消息到重试交换机,然后对当前消息做basicAck确认。这么做是因为我用的RabbitMQ版本对消息的deadLetterRoutingKey做了限制,直接nack且requeue=false只能进固定的死信队列,没法灵活控制“哪些消息进重试队列,哪些进死信队列”。通过手动重新投递,重试逻辑的控制权完全在自己手里,想要重试几次、间隔多久、header带什么信息,全都可控。
但是这种“发送新消息再确认老消息”的方式也有个坑:如果发送重试消息成功之后、basicAck之前,消费者进程崩了,那条消息会被重新投递,同时重试队列里也已经有一条新消息,两者会重复处理。因此消费端的幂等设计在这里尤其关键,一定要保证即使消息重复进入,最终对下游的数据写入是幂等无副作用的。
另一个细节是TTL的设置。我在重试消息上通过setExpiration("30000")设置了30秒的过期时间,这条消息进入重试队列后会被延迟30秒再路由回主队列。这里有个RabbitMQ的经典问题:RabbitMQ只会检查队列头部消息的TTL,也就是说如果有多条消息堆积在重试队列里,前面的消息没过期,后面的消息即使到了时间也不能被投递,这会拉长整体重试延迟。我在实践中发现这个行为特性确实存在,所以重试队列的消费能力要足够快,不要让重试消息大量堆积在重试队列里。
4.3 批量消费场景下的重试优化
大数据场景下,消费端很少一条一条处理消息,基本都是批量拉取。批量模式下重试机制的设计又有不同。假设一次拉取500条消息,批量写入下游数仓,因为下游某个字段类型不匹配导致整个批写入失败。这时候如果把500条消息全部requeue,大概率下次还是会一起失败,浪费大量资源。
我在处理批量失败时采用的策略是:先把这批消息缓存在本地内存里,逐条解析,找出真正的坏数据。解析失败的单条单独走重试逻辑,其余能正常处理的继续写入。这种方式能避免“一粒老鼠屎坏了一锅粥”,但要注意内存管理,本地缓存需要设置上限,防止大量消息堆积撑爆内存。
批量消费还有一个更务实的做法:对批量消息统一确认,但只对失败的子集做重试。如果这批消息里只有几条失败,就把成功的部分确认掉,失败的部分通过basicNack重新投递或者走重试队列。这个逻辑需要维护一个失败消息列表,在批量处理结果返回后统一处理。注意basicNack的multiple参数为true时是批量拒绝,使用时要看清语义,别把不该重试的消息也拒了。
4.4 高吞吐与重试的平衡调优
重试机制做重了,吞吐量必然会下降,因为每条失败消息都要额外走一次网络往返。我在压测中发现,如果消息失败率在1%左右,重试对吞吐量的影响还不明显;但失败率到5%以上,整个管道的吞吐量直接下降20%以上。原因是重试消息和正常消息混在同一个队列里,重试消息反复投递,消费端需要反复反序列化、解析、重试,占用大量CPU。
要平衡吞吐和可靠性,我的经验是把重试逻辑尽量从主链路上剥离开。比如失败的消息先投递到一个独立的“失败缓冲队列”,由单独的消费者处理重试策略,主消费者只负责正常业务。这样重试的并发度、资源占用和正常业务完全隔离,主链路吞吐不会因为少数失败消息而受拖累。
另一个调优方向是消费端的并发参数。默认情况下simple listener的并发消费者数是1,在大数据场景下几乎不可能满足吞吐要求。我一般会设置concurrency为10到20,max-concurrency为30,同时把prefetch值调大(比如100),让消费者一次拉取更多消息减少网络往返。但prefetch也不能太大,否则消费端内存压力大,而且消息长时间不被确认,Broker端重试投递的语义也会变得模糊。机器内存充足时prefetch=300比较均衡,内存紧张就回到100。
5. 常见问题与排查技巧实录
5.1 消息重复消费,怎么定位和解决
重复消费是重试机制最常被问到的问题。现象很好认:下游数据库里出现重复的订单记录,或者数仓某张表的行数比预期多。
排查思路分三步。首先看消费端的ack时机,确认消息是否在业务数据处理完成之前就被确认了。这种“先ack后处理”的写法我见过不少,一旦处理逻辑抛异常,消息已经确认无法重投,但业务也失败了,只能算“丢失”而不是“重复”。所以规范应该是先处理成功再ack,处理失败不ack。其次看异常场景下是否有线程安全问题导致重复处理,比如并发消费者数量大于1时,同一个业务ID的消息被两个线程同时消费。最后看重试投递逻辑,手动投递重试消息时是否把原始消息也保留了,导致同一业务数据有两条不同消息。
解决重复消费最根本的方法还是幂等。我在代码里要求所有数据消费逻辑必须对业务ID做唯一性检查,可以在数据库层面加唯一索引,也可以在Redis里用SETNX做前置判断。做大数据统计时尤其要注意幂等,因为聚合计算一旦重复执行,最终结果不正确还很难排查。
5.2 重试风暴:频繁重试打垮下游
重试风暴是我见过最具破坏性的故障模式。某天下游ES集群负载升高,响应变慢,上游大批消息消费失败,全员进入重试逻辑。重试消息投递回主队列后,消费端再次拉取再次失败,同时又产生新的重试消息,整个队列的消息数量指数级增长,下游被反复冲击,直到彻底不可用。
这种场景有两种有效应对手段。第一种是退避策略要带上限,重试间隔必须用指数退避,且设置最大间隔,防止重试频率过高。我一般会把最大重试间隔控制在60秒以上,即使某个下游故障持续几分钟,消费端的重试频率也不会对下游造成二次伤害。第二种是熔断机制,当下游连续失败率达到阈值,比如连续100条消息全部失败,主动暂停消费一段时间,让下游喘息恢复。
实现熔断最简单的方式是用一个计数器加定时器,在消费失败时递增计数,成功时清零。计数超过阈值就调用channel.basicCancel暂停消费,同时启动一个定时任务,30秒后重新恢复消费。这个逻辑虽然粗暴,但在生产环境实测非常有效,能快速止血,防止故障蔓延到整个集群。
5.3 消息一直不被消费,排查五大类原因
消息堆积但不消费,排查方向通常集中在五个方面。第一看消费者是否正常运行,是不是进程挂掉或者断开了连接。RabbitMQ管理界面可以看到Connection和Channel状态,消费者如果异常退出,队列的消费者数会变成0。第二看消费端是否设置了concurrency=0,这种情况消费者根本不会启动。第三看消息是否被requeue策略卡住,有些时候消费端一直在抛异常但异常被吞掉,消息被无限重试,从外部看就是一直在消费但又一直不成功。第四看死信队列里是否有消息堆积,重试次数耗尽的消息全去了死信,主队列看似正常,实际业务早已停摆。第五看prefetch和消费者的处理速度是否匹配,如果单条消息处理耗时长而prefetch很小,队列里的消息消耗速度就会很慢。
排查时我最常用的手段是查看队列的Unrouted消息数和各队列的Ready/Unacked数量对比。Ready一直涨说明消费者拉取不过来,Unacked一直很高说明消费者拿到消息后处理很慢或者卡住了。配合管理界面的Message Rates图表,很快能定位到瓶颈。
5.4 手动投递重试消息时的连接与内存坑
手动投递重试消息时有一个容易踩的坑:RabbitTemplate如果使用不当会造成连接泄漏。默认情况下RabbitTemplate每次发送都会从连接工厂获取一个连接,高频发送时如果连接池配置不合理,会出现连接数飙升的问题。我建议在配置文件里明确指定缓存连接模式,并且在生产环境使用CachingConnectionFactory,把channelCacheSize设成一个合适的值,比如50,避免频繁创建通道的开销。
另一个坑是Message对象的复用。在批量失败场景中,如果循环发送多条重试消息,需要为每条消息创建独立的MessageProperties,避免多个消息共享同一个MessageProperties实例导致header覆盖。我早期写循环发送时就踩过这个坑,所有消息的retryCount头都变成了同一个值,重试上限判断完全失效。
内存方面的坑主要集中在TTL和堆积上。重试队列如果TTL设置很长,比如5分钟,而失败消息量大,重试队列本身就会堆积大量延迟消息,消耗大量内存。我通常建议把延迟阈值控制在1分钟内,配合死信队列做兜底处理。如果业务上确实需要更长延迟,优先考虑RabbitMQ的延迟插件(rabbitmq-delayed-message-exchange),而不是单纯依靠TTL实现,因为延迟插件的调度机制比TTL扫队列高效得多。
5.5 运维侧的经验:监控告警配置清单
最后说说运维侧必须盯的几个指标。队列深度是最基本的,主要队列和重试队列、死信队列都要单独监控。我习惯在Grafana里配置三张看板:主队列深度趋势、死信队列入队速率、消费者处理延迟。死信队列入队速率这个指标特别重要,一旦出现持续上涨,基本可以断定有系统性问题在发生。
消息处理耗时同样需要监控。RabbitMQ管理界面提供了消费者处理时间的分布,如果平均值持续走高,说明业务处理逻辑有瓶颈需要优化。还有一个容易被忽视的指标是channel的关闭频率,如果大量channel频繁创建和关闭,通常是消费者的连接稳定性出了问题。
我给这套监控体系配置了两级告警。第一级是死信队列有消息进入就报警,用企业微信或者钉钉机器人推送,提醒值班人员关注。第二级是队列深度超过阈值,比如主队列超过10万、死信队列超过5000时告警。阈值不能设得太小,否则正常的流量波动都会触发告警产生疲劳;也不能太大,否则故障半小时才发现,损失已经造成。这个平衡需要根据各自业务的实际吞吐量来调整,不能照搬别人的参数。
6. 最后分享几个实战中沉淀下来的心得
代码写到这里,最后聊几个我实操过程中比较深刻的体会。
第一,重试机制不是越复杂越好。早期我做过一套带有重试状态机的方案,消息在多个队列之间流转,每一步都有状态记录。问题是出问题时定位太困难了,一条消息在哪个队列、处于哪个阶段、接下来要去哪,全靠日志串联,排查效率极低。后来我简化成“主队列+重试队列+死信队列”的三段式,出问题直接看死信队列就能找到答案。可靠性和可维护性权衡下来,简单方案更适合大多数团队。
第二,不要迷信框架的自动重试。Spring Retry确实写起来方便,但它隐藏了消息的状态流转,生产环境一旦遇到问题,CRITICAL级别的故障排查会把时间浪费在“这条消息到底被重试了几次”这类问题上。手动ack+显式重试虽然代码量多一些,但每一步逻辑都透明可见,对大数据链路这种对可靠性要求极高的场景,透明比简单更重要。
第三,任何重试方案都必须在压测环境验证过。我见过一个项目,重试方案在测试环境一切正常,上线后才发现高并发下重试消息的TTL过期时间被大量消息撑爆,延迟从预期的30秒变成了10分钟。压测时一定要模拟下游故障的场景,观察重试风暴、队列堆积、消费恢复这些关键节点,数据链路无小事,每条消息背后都是一份不能丢的数据。