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
相关产品推荐
相关产品推荐

