迁移至AWS SDK v2时,如何重写S3 CreateEvent的SQS监听器?
AWS SDK v2 迁移:Spring Boot 中通过 @SqsListener 处理 S3 事件
AWS SDK v2 确实移除了 S3EventNotification 类,迁移时需要直接解析SQS消息的JSON内容来实现原有逻辑,具体步骤如下:
1. 调整方法参数类型
原来的方法参数是S3EventNotification,现在改为接收原始的消息字符串(或Message<String>以获取消息元数据)——因为S3发送到SQS的事件本质是JSON格式的字符串。
2. 解析S3事件JSON
可以用Jackson的ObjectMapper来解析JSON,有两种常用实现方式:
方式一:自定义实体类映射(类型安全,推荐)
先定义和S3事件结构匹配的实体类,模仿原S3EventNotification的结构:
import com.fasterxml.jackson.annotation.JsonProperty; import java.util.List; public class S3EventV2 { @JsonProperty("Records") private List<S3EventRecordV2> records; public List<S3EventRecordV2> getRecords() { return records; } public void setRecords(List<S3EventRecordV2> records) { this.records = records; } public static class S3EventRecordV2 { @JsonProperty("eventName") private String eventName; // 根据业务需求,可添加其他字段(如S3对象信息、桶名等) public String getEventName() { return eventName; } public void setEventName(String eventName) { this.eventName = eventName; } } }
然后修改监听方法:
import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.cloud.aws.messaging.annotation.SqsListener; import org.springframework.cloud.aws.messaging.listener.SqsMessageDeletionPolicy; import java.io.IOException; @SqsListener(value = "[the physical id]", deletionPolicy = SqsMessageDeletionPolicy.ON_SUCCESS) public void handleMessage(final String sqsMessageBody) throws IOException { // 建议注入Spring容器中的ObjectMapper实例,而非每次新建 ObjectMapper objectMapper = new ObjectMapper(); S3EventV2 s3Event = objectMapper.readValue(sqsMessageBody, S3EventV2.class); // 增加非空判断,避免空指针异常 if (!s3Event.getRecords().isEmpty() && s3Event.getRecords().get(0).getEventName().startsWith("ObjectCreated")) { // 执行原有的S3对象创建事件处理逻辑 [process the S3 createObject event] } }
方式二:直接用Map解析(快速实现,无需实体类)
如果不需要类型安全,也可以直接将JSON解析为Map:
import com.fasterxml.jackson.databind.ObjectMapper; import org.springframework.cloud.aws.messaging.annotation.SqsListener; import org.springframework.cloud.aws.messaging.listener.SqsMessageDeletionPolicy; import java.io.IOException; import java.util.List; import java.util.Map; @SqsListener(value = "[the physical id]", deletionPolicy = SqsMessageDeletionPolicy.ON_SUCCESS) public void handleMessage(final String sqsMessageBody) throws IOException { ObjectMapper objectMapper = new ObjectMapper(); Map<String, Object> eventMap = objectMapper.readValue(sqsMessageBody, Map.class); List<Map<String, Object>> records = (List<Map<String, Object>>) eventMap.get("Records"); if (records != null && !records.isEmpty()) { String eventName = (String) records.get(0).get("eventName"); if (eventName != null && eventName.startsWith("ObjectCreated")) { // 执行处理逻辑 [process the S3 createObject event] } } }
注意事项
- 确保依赖正确:Spring Cloud AWS v2.4.4需要搭配AWS SDK v2的相关依赖(如
software.amazon.awssdk:s3、org.springframework.cloud:spring-cloud-starter-aws-messaging) - 避免空指针:必须增加对records列表、eventName的非空判断,原代码直接调用
get(0)存在空指针风险 - 复用ObjectMapper:建议在Spring容器中配置并注入
ObjectMapper实例,不要每次方法调用时新建,提升性能
内容的提问来源于stack exchange,提问作者Martin Bamford
相关产品推荐
相关产品推荐

