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

使用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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.10.06 09:45:02