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

Spring SQS:发送至死信队列前添加自定义消息属性遇异常求助

解决方案

1. 修复切面的UnsupportedOperationException问题

你遇到的异常是因为SendMessageRequest.getMessageAttributes()返回的是不可修改的Map(AWS SDK内部用Collections.unmodifiableMap包装),直接调用put会触发不支持操作异常。修改切面代码,先将原属性复制到可变Map中再添加自定义属性:

@Before(value = "execution(* com.amazonaws.services.sqs.AmazonSQS*Client.sendMessage*(com.amazonaws.services.sqs.model.SendMessageRequest,..)) && args(request,..)",
        argNames = "request")
public void before(SendMessageRequest request) {
    Map<String, MessageAttributeValue> originalAttributes = request.getMessageAttributes();
    // 处理null或不可变Map的情况,转成可变的HashMap
    Map<String, MessageAttributeValue> mutableAttributes = new HashMap<>();
    if (originalAttributes != null) {
        mutableAttributes.putAll(originalAttributes);
    }
    // 添加自定义死信标识属性
    mutableAttributes.put("MoveToDlqFlag", createMessageAttribute("true"));
    // 将可变Map设置回请求对象
    request.setMessageAttributes(mutableAttributes);
}

private MessageAttributeValue createMessageAttribute(String value) {
    MessageAttributeValue attr = new MessageAttributeValue();
    attr.setDataType("String");
    attr.setStringValue(value);
    return attr;
}

2. Spring Cloud AWS 原生死信增强方案

如果不想用切面,Spring Cloud AWS SQS提供了更贴合框架的死信处理方式:

方式一:自定义DeadLetterQueueResolver

实现DeadLetterQueueResolver接口,在解析死信队列时修改消息属性并手动发送:

@Component
public class CustomDeadLetterQueueResolver implements DeadLetterQueueResolver {
    private final AmazonSQS amazonSQS;

    public CustomDeadLetterQueueResolver(AmazonSQS amazonSQS) {
        this.amazonSQS = amazonSQS;
    }

    @Override
    public Queue getDeadLetterQueue(Queue sourceQueue, Message message) {
        // 从源队列属性中解析死信队列ARN
        String redrivePolicy = sourceQueue.getAttributes().get("RedrivePolicy");
        String dlqArn = extractDlqArnFromRedrivePolicy(redrivePolicy);
        Queue dlq = Queue.fromArn(dlqArn);
        
        // 构造死信消息,复制原属性并添加自定义字段
        SendMessageRequest dlqRequest = new SendMessageRequest()
                .withQueueUrl(dlq.getUrl())
                .withMessageBody(message.getBody());
        
        Map<String, MessageAttributeValue> mutableAttrs = new HashMap<>(message.getMessageAttributes());
        mutableAttrs.put("MoveToDlqFlag", createMessageAttribute("true"));
        mutableAttrs.put("FailureTime", createMessageAttribute(String.valueOf(System.currentTimeMillis())));
        dlqRequest.setMessageAttributes(mutableAttrs);
        
        // 手动发送至死信队列
        amazonSQS.sendMessage(dlqRequest);
        
        // 返回null阻止框架自动发送死信
        return null;
    }

    private String extractDlqArnFromRedrivePolicy(String redrivePolicy) {
        ObjectMapper mapper = new ObjectMapper();
        try {
            JsonNode node = mapper.readTree(redrivePolicy);
            return node.get("deadLetterTargetArn").asText();
        } catch (IOException e) {
            throw new RuntimeException("解析死信队列配置失败", e);
        }
    }

    private MessageAttributeValue createMessageAttribute(String value) {
        MessageAttributeValue attr = new MessageAttributeValue();
        attr.setDataType("String");
        attr.setStringValue(value);
        return attr;
    }
}

方式二:监听方法内手动处理死信

关闭队列原生的RedrivePolicy,在监听方法中捕获异常,手动构造死信消息并发送:

@SqsListener(value = {"${sqs.queue.source-queue}"}, deletionPolicy = ON_SUCCESS)
public void handleMessageReceived(String rawMessage, @Header("SenderId") String senderId, @Headers Map<String, Object> header) {
    try {
        var orderMessageWrapper = messageWrapperUtil.create(rawMessage, Order.class);
        var order = orderMessageWrapper.getMessage();
        var receiveCount = (Integer) header.get("ApproximateReceiveCount");
        
        // 业务逻辑处理
        ...
    } catch (Exception e) {
        // 构造死信请求
        SendMessageRequest dlqRequest = new SendMessageRequest()
                .withQueueUrl("${sqs.queue.dlq-url}")
                .withMessageBody(rawMessage);
        
        // 添加自定义属性
        Map<String, MessageAttributeValue> attrs = new HashMap<>();
        attrs.put("MoveToDlqFlag", createMessageAttribute("true"));
        attrs.put("FailureReason", createMessageAttribute(e.getMessage()));
        attrs.put("ReceiveCount", createMessageAttribute(receiveCount.toString()));
        dlqRequest.setMessageAttributes(attrs);
        
        // 发送死信并删除原消息
        amazonSQS.sendMessage(dlqRequest);
        amazonSQS.deleteMessage(new DeleteMessageRequest("${sqs.queue.source-queue}", (String) header.get("ReceiptHandle")));
    }
}

private MessageAttributeValue createMessageAttribute(String value) {
    MessageAttributeValue attr = new MessageAttributeValue();
    attr.setDataType("String");
    attr.setStringValue(value);
    return attr;
}

3. 注意事项

  • 切面方式要确保能拦截到框架发送死信的sendMessage调用(Spring Cloud AWS会在重试耗尽后调用SDK发送死信)。
  • 手动处理死信时,必须关闭队列原生的RedrivePolicy,避免重复发送。
  • 自定义属性需符合SQS限制:单个键长度≤256字符,值≤1KB,总属性大小≤256KB。

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.08.19 12:31:26