使用Java SDK上传Kinesis数据流到S3桶报InvalidRedirectLocation错误如何解决
报错原因
你当前调用的PutObjectRequest重载方法的第三个入参是S3静态网站托管的重定向地址,而非要上传的文件内容,你传入的Kinesis ARN不符合重定向地址必须以http:///https:////开头的规则,因此触发InvalidRedirectLocation错误。
正确实现方案
要实现关闭Kinesis数据流时将流内数据上传到S3,需要先读取Kinesis流中的实际记录,再将数据写入S3,完整实现逻辑如下:
前置依赖
你的项目中需要引入AWS Java SDK的Kinesis和S3相关依赖,以Maven为例:
<dependency> <groupId>com.amazonaws</groupId> <artifactId>aws-java-sdk-s3</artifactId> <version>1.12.XXX</version> </dependency> <dependency> <groupId>com.amazonaws</groupId> <artifactId>aws-java-sdk-kinesis</artifactId> <version>1.12.XXX</version> </dependency>
替换版本号为你项目使用的AWS SDK对应版本即可。
完整实现代码
import com.amazonaws.services.kinesis.AmazonKinesis; import com.amazonaws.services.kinesis.AmazonKinesisClientBuilder; import com.amazonaws.services.kinesis.model.*; import com.amazonaws.services.s3.AmazonS3; import com.amazonaws.services.s3.AmazonS3ClientBuilder; import com.amazonaws.services.s3.model.PutObjectRequest; import java.io.ByteArrayInputStream; import java.nio.charset.StandardCharsets; import java.util.List; import java.util.stream.Collectors; public void uploadClosedKinesisStreamToS3() { // 初始化客户端 AmazonKinesis kinesisClient = AmazonKinesisClientBuilder.standard() .withRegion(awsRegion) .build(); AmazonS3 s3Client = AmazonS3ClientBuilder.standard() .withRegion(awsRegion) .build(); String streamName = connectionRequestRepository.findStream(); String bucketName = "downloadable-cases"; String fileName = streamName + ".json"; try { // 1. 先校验Kinesis流状态,确认已关闭/停止写入,可根据你的业务判定规则调整状态判断逻辑 StreamDescription streamDesc = kinesisClient.describeStream(streamName).getStreamDescription(); if (!"ACTIVE".equals(streamDesc.getStreamStatus())) { // 2. 读取流内所有分片的记录 StringBuilder allRecordsContent = new StringBuilder(); List<Shard> shards = streamDesc.getShards(); for (Shard shard : shards) { String shardIterator = kinesisClient.getShardIterator(new GetShardIteratorRequest() .withStreamName(streamName) .withShardId(shard.getShardId()) .withShardIteratorType(ShardIteratorType.TRIM_HORIZON)).getShardIterator(); GetRecordsResult recordsResult; do { recordsResult = kinesisClient.getRecords(new GetRecordsRequest().withShardIterator(shardIterator)); // 把记录数据转成字符串拼接,可根据你的业务需要调整格式 List<String> recordData = recordsResult.getRecords().stream() .map(record -> new String(record.getData().array(), StandardCharsets.UTF_8)) .collect(Collectors.toList()); recordData.forEach(allRecordsContent::append); shardIterator = recordsResult.getNextShardIterator(); } while (shardIterator != null && recordsResult.getRecords().size() > 0); } // 3. 将拼接好的流数据上传到S3 byte[] contentBytes = allRecordsContent.toString().getBytes(StandardCharsets.UTF_8); ByteArrayInputStream inputStream = new ByteArrayInputStream(contentBytes); s3Client.putObject(new PutObjectRequest(bucketName, fileName, inputStream, null)); } } catch (Exception e) { e.printStackTrace(); } finally { kinesisClient.shutdown(); s3Client.shutdown(); } }
特殊场景说明
如果你只是需要把Kinesis流的ARN作为纯文本内容存储到S3,不需要读取流内实际数据,仅需要修改你原有代码的传参方式即可:
public void uploadStreamArnToS3() { AmazonS3 s3Client = AmazonS3ClientBuilder.standard() .withRegion(String.valueOf(awsRegion)) .build(); try { String fileName = connectionRequestRepository.findStream() +".json"; String bucketName = "downloadable-cases"; String locationData = "arn:aws-us-***-1:kinesis:***:stream/" + connectionRequestRepository.findStream(); // 把ARN字符串转成输入流再传入PutObjectRequest即可避免重定向报错 ByteArrayInputStream inputStream = new ByteArrayInputStream(locationData.getBytes(StandardCharsets.UTF_8)); s3Client.putObject(new PutObjectRequest(bucketName, fileName, inputStream, null)); } catch (AmazonServiceException ex) { System.out.println("Error: " + ex.getMessage()); } }
内容的提问来源于stack exchange,提问作者Tip
相关产品推荐
相关产品推荐

