You need to enable JavaScript to run this app.
优惠活动
大模型
产品
解决方案
定价
更多

RabbitMQ重试机制下仅首次保存消息至数据库的方案问询

Fix Duplicate Database Saves During RabbitMQ Retries

Got it, let's tackle this problem where your RabbitMQ messages are being saved to the database multiple times during retries. The core issue here is that your messageRepository.save() call lives inside the listener method, which gets re-executed every time the retry interceptor kicks in. Here are a few clean solutions to fix this:

Solution 1: Check Retry Count in the Listener (Quick Fix)

Use Spring Retry's RetrySynchronizationManager to get the current retry attempt count, and only save the message on the first try. This requires minimal changes to your existing listener code:

@Slf4j
@Component
@AllArgsConstructor
public class MessageListener {
    private final MessageRepository messageRepository;
    private final MessageService messageService;

    @RabbitListener(queues = MessageConfiguration.QUEUE, containerFactory = MessageConfiguration.FACTORY_CONTAINER_NAME)
    public void process(SomeMessage someMessage) {
        RetryContext retryContext = RetrySynchronizationManager.getContext();
        // Only save on the first attempt (retryCount starts at 0 for initial execution)
        if (retryContext == null || retryContext.getRetryCount() == 0) {
            messageRepository.save(someMessage); // Transform to entity first if needed
        }
        messageService.process(someMessage);
    }
}

How it works:

  • RetrySynchronizationManager binds a RetryContext to the current thread when the retry interceptor runs.
  • For the first execution, retryCount is 0. Each subsequent retry increments this count.
  • If the message succeeds on the first try, retryContext might be null (depends on interceptor setup), so we handle that case too.

Solution 2: Decouple Save Logic with a Pre-Retry Advice (Cleaner Separation)

Move the save logic into a custom MethodBeforeAdvice that runs before the retry interceptor. This keeps your listener focused solely on business processing, following the single responsibility principle.

First, create the advice component:

@Component
public class FirstAttemptSaveAdvice implements MethodBeforeAdvice {
    private final MessageRepository messageRepository;

    public FirstAttemptSaveAdvice(MessageRepository messageRepository) {
        this.messageRepository = messageRepository;
    }

    @Override
    public void before(Method method, Object[] args, Object target) throws Throwable {
        RetryContext retryContext = RetrySynchronizationManager.getContext();
        // Only save if this is the first attempt
        if (retryContext == null || retryContext.getRetryCount() == 0) {
            // Extract the message from method arguments (adjust index if your listener has different params)
            if (args.length > 0 && args[0] instanceof SomeMessage) {
                SomeMessage someMessage = (SomeMessage) args[0];
                messageRepository.save(someMessage); // Transform to entity first if needed
            }
        }
    }
}

Then update your container factory to include this advice before the retry interceptor (advice chain order matters!):

@Bean(name = FACTORY_CONTAINER_NAME)
public SimpleRabbitListenerContainerFactory factoryQueueExample(ConnectionFactory connectionFactory,
                                                               RetryOperationsInterceptor defaultRetryOperationsInterceptor,
                                                               FirstAttemptSaveAdvice firstAttemptSaveAdvice) {
    SimpleRabbitListenerContainerFactory factory = new SimpleRabbitListenerContainerFactory();
    factory.setConnectionFactory(connectionFactory);
    factory.setMessageConverter(getMessageConverter());
    factory.setDefaultRequeueRejected(false);
    // Execute save advice first, then retry interceptor
    Advice[] adviceChain = {firstAttemptSaveAdvice, defaultRetryOperationsInterceptor};
    factory.setAdviceChain(adviceChain);
    return factory;
}

Finally, simplify your listener by removing the save call:

@Slf4j
@Component
@AllArgsConstructor
public class MessageListener {
    private final MessageService messageService;

    @RabbitListener(queues = MessageConfiguration.QUEUE, containerFactory = MessageConfiguration.FACTORY_CONTAINER_NAME)
    public void process(SomeMessage someMessage) {
        messageService.process(someMessage);
    }
}

Solution 3: Save with Retry Attempt Number (For Audit Logs)

If you need to track all retry attempts (instead of just the first), modify your save logic to include the attempt number. This is useful for debugging and audit purposes:

@Slf4j
@Component
@AllArgsConstructor
public class MessageListener {
    private final MessageRepository messageRepository;
    private final MessageService messageService;

    @RabbitListener(queues = MessageConfiguration.QUEUE, containerFactory = MessageConfiguration.FACTORY_CONTAINER_NAME)
    public void process(SomeMessage someMessage) {
        RetryContext retryContext = RetrySynchronizationManager.getContext();
        // Calculate attempt number (1 = first try, 2 = first retry, etc.)
        int attemptNumber = retryContext != null ? retryContext.getRetryCount() + 1 : 1;
        
        // Convert message to entity and add attempt number
        MessageEntity messageEntity = convertToEntity(someMessage);
        messageEntity.setAttemptNumber(attemptNumber);
        
        messageRepository.save(messageEntity);
        messageService.process(someMessage);
    }

    // Helper method to convert SomeMessage to your database entity
    private MessageEntity convertToEntity(SomeMessage someMessage) {
        // Implement your conversion logic here
        return new MessageEntity();
    }
}

Key Notes:

  • All solutions rely on Spring Retry's thread-bound RetryContext, which is safe for multi-threaded listener containers.
  • For stateless retries (your current setup), the RetryContext is reused across retries for the same message, so the retryCount increments correctly.

内容的提问来源于stack exchange,提问作者Dherik

相关产品推荐
方舟 Agent Plan

超全模态模型 × Harness 升级,最新支持 Deepseek-V4.1-Flash、GLM-5.3 系列、Doubao-Seedream-5.0-pro、Kimi-K3 (部分), 限时 9.9 元起

最近更新时间:2026.05.27 07:28:03