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

@SqsListener反序列化失败时如何删除SQS队列中的对应消息?

问题原因

你当前使用的SqsMessageDeletionPolicy.ON_SUCCESS删除策略,仅在监听方法正常执行完成、无任何异常抛出时才会删除对应消息。而MessageConversionException属于消息转换阶段的异常,发生在SQS原始消息反序列化为User对象的过程中,此时还未进入你编写的receiveEvent业务方法,方法上声明的@MessageExceptionHandler无法捕获该异常,异常向上抛出后会触发SQS默认重试机制,导致消息反复被消费。

可行解决方案
  • 方案1:全局配置拦截器统一处理消息转换异常
    自定义SQS监听容器工厂,添加全局消息拦截器,提前捕获MessageConversionException,异常被拦截后不会向上抛出,即可触发自动删除规则:
import org.springframework.context.annotation.Bean;
import org.springframework.context.annotation.Configuration;
import software.amazon.awssdk.services.sqs.SqsAsyncClient;
import org.springframework.cloud.aws.sqs.config.SqsMessageListenerContainerFactory;
import org.springframework.messaging.converter.MessageConversionException;
import java.util.concurrent.CompletableFuture;

@Configuration
public class SqsGlobalConfig {

    @Bean
    public SqsMessageListenerContainerFactory<Object> defaultSqsListenerContainerFactory(
            SqsAsyncClient sqsAsyncClient) {
        return SqsMessageListenerContainerFactory
                .builder()
                .sqsAsyncClient(sqsAsyncClient)
                .configure(options -> options
                        .messageListenerInterceptor((message, executionChain) -> {
                            try {
                                return executionChain.execute(message);
                            } catch (MessageConversionException e) {
                                // 可在此处添加日志,记录坏消息内容用于排查问题
                                return CompletableFuture.completedFuture(null);
                            }
                        })
                )
                .build();
    }
}

该方案适合所有SQS监听接口都需要统一处理反序列化异常的场景,无需修改业务代码。

  • 方案2:接收原始消息手动反序列化
    修改监听方法入参为原始字符串类型,在业务方法内自行完成反序列化逻辑,反序列化异常完全由业务代码控制:
import com.fasterxml.jackson.core.JsonProcessingException;
import com.fasterxml.jackson.databind.ObjectMapper;
import io.awspring.cloud.sqs.annotation.SqsListener;
import io.awspring.cloud.sqs.listener.SqsMessageDeletionPolicy;

@SqsListener(value = "${cloud.aws.sqs.url}", deletionPolicy = SqsMessageDeletionPolicy.ON_SUCCESS)
public void receiveEvent(String rawMessage) {
    User user;
    try {
        user = new ObjectMapper().readValue(rawMessage, User.class);
    } catch (JsonProcessingException e) {
        // 反序列化失败时直接返回,不抛出异常即可触发消息删除
        // 可在此处添加日志记录坏消息内容
        return;
    }
    handleRequest(user);
}

该方案无需修改全局配置,灵活度更高,适合仅单个监听接口需要处理反序列化异常的场景。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.30 23:06:04