Spring Boot中RabbitMQ队列监听器执行逻辑不一致问题咨询
问题根因
该问题和RabbitMQ的分布式特性无关,是由以下几个原因共同导致的:
- 无状态重试逻辑缺陷:你配置的是
stateless无状态重试,只要整个监听方法抛出未捕获的异常,就会从头重新执行整个方法,不会记录上一次的执行进度。如果某次重试中前半部分操作执行成功,后半部分抛出异常,下次重试时哪怕createNeoPost因为数据已变更抛出异常,之前已经执行成功的incrementMediaCount也不会回滚,就会出现你看到的“create抛异常但计数仍增加”的现象。 - 无事务边界+操作非幂等:三个数据操作分别操作Neo4j和PostgreSQL两个数据源,没有配置统一事务,且没有做幂等校验,每一步操作成功后都会直接持久化,不会因为后续操作失败而回滚,重试多次就会导致计数被重复累加。
- 日志存储逻辑位置靠后:日志保存是监听方法的最后一步,只要前面任意一步抛出异常终止执行,日志就不会写入,所以最终存储的日志数量远小于消息实际接收次数。
- 额外排查点:可以先确认
ResourceNotFoundException是否被全局异常处理器捕获后没有重新抛出,如果异常没有透传到Spring Rabbit的重试拦截器层面,也会导致前面抛异常后后面代码继续执行,且不会触发重试。
修复方案
1. 快速修复(适合不需要严格强一致的场景)
首先给所有业务操作加幂等校验,避免重复执行:
// NeoPostService新增幂等校验方法 fun existsByPostId(postId: Long): Boolean { return neoPostRepository.existsByPostId(postId) } // 调整监听方法逻辑 @RabbitListener(queues = [RabbitConfiguration.Companion.Queues.POST_CREATION_QUEUE], concurrency = "3") fun postCreationListener(postMessage: PostMessage) { // 幂等校验:已经处理过的消息直接跳过 if (neoPostService.existsByPostId(postMessage.postId)) { return } neoPostService.createNeoPost(postMessage.postId) accountService.incrementMediaCount(postMessage.userId) loggingRepository.save( Log( resourceType = ResourceType.Post, entityEventType = EntityEventType.Created, message = "Post ${postMessage.postId} created for user ${postMessage.userId}" ) ) }
同时把重试拦截器调整为有状态重试,避免不必要的全方法重跑:
@Bean fun retryInterceptor(): RetryOperationsInterceptor? { return RetryInterceptorBuilder.stateful() .backOffOptions(100, 5.0, 10000) .maxAttempts(5) .recoverer(RejectAndDontRequeueRecoverer()) .build() }
2. 严格一致方案(适合对数据准确性要求高的场景)
引入JTA分布式事务管控跨数据源操作,或者改用最终一致性方案:
- 先把消息ID和处理状态存入操作流水表,再执行业务逻辑,所有操作完成后更新流水状态为成功
- 新增定时任务定时扫描超过一定时间未处理成功的流水,执行补偿操作
这样哪怕重试多次,也不会出现数据不一致的问题。
内容的提问来源于stack exchange,提问作者Annon
相关产品推荐
相关产品推荐

