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

迁移至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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.06 09:15:13