AWS SDK中S3EventNotification位置及SQS消息序列化咨询
问题
以下是我的SQS监听器类:
@RequiredArgsConstructor @Slf4j @Service public class SQSS3CreatedNotificationService { private final FluxProperties fluxProperties; private final ConsolidateDocumentInputPort consolidateDocumentInputPort; @SqsListener("#{fluxProperties.getS3CreatedNotificationQueueName()}") public void listen(final Message s3EventNotification) { log.trace("Message received: {}", s3EventNotification.body()); consolidateDocumentInputPort.consolidate(UUID.randomUUID()); } }
如您所见,我使用Message类参数接收SQS消息,该消息是一个S3事件。以下是log.trace("Message received: {}", s3EventNotification.body());的格式化输出:
{ "Records": [ { "eventVersion": "2.1", "eventSource": "aws:s3", "awsRegion": "us-east-1", "eventTime": "2024-02-16T13:35:26.982Z", "eventName": "ObjectCreated:Put", "userIdentity": { "principalId": "AIDAJDPLRKLG7UEXAMPLE" }, "requestParameters": { "sourceIPAddress": "127.0.0.1" }, "responseElements": { "x-amz-request-id": "234d0f2e", "x-amz-id-2": "eftixk72aD6Ap51TnqcoF8eFidJG9Z/2" }, "s3": { "s3SchemaVersion": "1.0", "configurationId": "38240f0e", "bucket": { "name": "espaidoc", "ownerIdentity": { "principalId": "A3NL1KOZZKExample" }, "arn": "arn:aws:s3:::espaidoc" }, "object": { "key": "references/reference1.sha", "sequencer": "0055AED6DCD90281E5", "size": 2333, "eTag": "1d11469e9d81f07729548d7708bbab82" } } } ] }
我想咨询:AWS SDK中S3EventNotification类的位置在哪里?能否将该SQS消息序列化为SDK中已定义的该类?
回答
1. S3EventNotification类的位置
该类属于AWS Java SDK v1的S3模块,完整包路径为:com.amazonaws.services.s3.event.S3EventNotification
如果使用构建工具管理依赖:
- Maven需引入
aws-java-sdk-s3依赖 - Gradle需引入
com.amazonaws:aws-java-sdk-s3依赖
2. 序列化SQS消息到该类
完全可以将SQS消息序列化为S3EventNotification类,有两种常见方式:
方式一:利用Spring Cloud AWS自动序列化
直接修改监听器方法的参数类型为S3EventNotification,Spring Cloud AWS会自动完成消息体的序列化:
@RequiredArgsConstructor @Slf4j @Service public class SQSS3CreatedNotificationService { private final FluxProperties fluxProperties; private final ConsolidateDocumentInputPort consolidateDocumentInputPort; @SqsListener("#{fluxProperties.getS3CreatedNotificationQueueName()}") public void listen(final S3EventNotification s3EventNotification) { log.trace("Message received: {}", s3EventNotification); // 提取事件详情示例 S3EventNotification.S3EventNotificationRecord record = s3EventNotification.getRecords().get(0); String bucketName = record.getS3().getBucket().getName(); String objectKey = record.getS3().getObject().getKey(); log.info("S3对象已创建:桶{},路径{}", bucketName, objectKey); consolidateDocumentInputPort.consolidate(UUID.randomUUID()); } }
方式二:手动序列化(保留Message参数)
如果需要保留Message参数,可以使用Jackson手动解析消息体:
@RequiredArgsConstructor @Slf4j @Service public class SQSS3CreatedNotificationService { private final FluxProperties fluxProperties; private final ConsolidateDocumentInputPort consolidateDocumentInputPort; private final ObjectMapper objectMapper; // 注入Spring默认的ObjectMapper @SqsListener("#{fluxProperties.getS3CreatedNotificationQueueName()}") public void listen(final Message s3EventNotification) throws IOException { log.trace("Message received: {}", s3EventNotification.body()); // 手动解析为S3EventNotification S3EventNotification notification = objectMapper.readValue(s3EventNotification.body(), S3EventNotification.class); consolidateDocumentInputPort.consolidate(UUID.randomUUID()); } }
内容的提问来源于stack exchange,提问作者Jordi
相关产品推荐
相关产品推荐

