<think>在处理从 `PipedInputStream` 读取数据且无法预先知道流大小的情况下,使用 AWS SDK for Java 2.40.0 的 `AsyncRequestBody` 上传到 S3 时,需要采取特殊措施。以下是具体解决方案和配置方法:
### 核心问题
- **无法提供 `contentLength`**:`AsyncRequestBody.fromInputStream()` 要求明确的内容长度,但动态生成的流(如通过管道传输)可能无法预先确定大小。
- **管道流限制**:`PipedInputStream` 默认缓冲区较小(1024字节),可能导致写入线程阻塞或读取端异常(如 `IOException: Read end dead`)[^1][^2][^3]。
### 解决方案
#### 步骤1:扩大管道缓冲区
调整 `PipedInputStream` 的缓冲区大小,避免读写线程阻塞:
```java
// 创建管道流时指定缓冲区大小(例如 1MB)
PipedOutputStream pout = new PipedOutputStream();
PipedInputStream pin = new PipedInputStream(pout, 1024 * 1024); // 1MB 缓冲区
```
#### 步骤2:使用分块上传(Multipart Upload)
当无法提供总大小时,改用 S3 的分块上传 API:
1. **初始化分块上传**:
```java
CreateMultipartUploadRequest createRequest = CreateMultipartUploadRequest.builder()
.bucket("your-bucket")
.key("object-key")
.build();
CreateMultipartUploadResponse response = s3AsyncClient.createMultipartUpload(createRequest).join();
String uploadId = response.uploadId();
```
2. **分块读取并上传**:
```java
List<CompletedPart> completedParts = new ArrayList<>();
byte[] buffer = new byte[5 * 1024 * 1024]; // 5MB 分块(S3 要求最小5MB)
int partNumber = 1;
int bytesRead;
while ((bytesRead = pin.read(buffer)) != -1) {
// 创建分片请求体(明确提供当前分片大小)
AsyncRequestBody partBody = AsyncRequestBody.fromBytes(buffer, 0, bytesRead);
UploadPartRequest uploadRequest = UploadPartRequest.builder()
.bucket("your-bucket")
.key("object-key")
.uploadId(uploadId)
.partNumber(partNumber)
.build();
// 上传分片
CompletedPart part = s3AsyncClient.uploadPart(uploadRequest, partBody)
.thenApply(uploadPartResponse ->
CompletedPart.builder()
.partNumber(partNumber)
.eTag(uploadPartResponse.eTag())
.build()
).join();
completedParts.add(part);
partNumber++;
}
```
3. **完成上传**:
```java
CompletedMultipartUpload completedUpload = CompletedMultipartUpload.builder()
.parts(completedParts)
.build();
CompleteMultipartUploadRequest completeRequest =
CompleteMultipartUploadRequest.builder()
.bucket("your-bucket")
.key("object-key")
.uploadId(uploadId)
.multipartUpload(completedUpload)
.build();
s3AsyncClient.completeMultipartUpload(completeRequest).join();
```
### 关键优化
1. **动态缓冲区管理**:
- 使用 `ByteArrayOutputStream` 临时缓存数据,直到达到分块大小(5MB):
```java
ByteArrayOutputStream chunkBuffer = new ByteArrayOutputStream(5 * 1024 * 1024);
byte[] tempBuffer = new byte[8192];
int bytesRead;
while ((bytesRead = pin.read(tempBuffer)) != -1) {
chunkBuffer.write(tempBuffer, 0, bytesRead);
if (chunkBuffer.size() >= 5 * 1024 * 1024) {
uploadChunk(chunkBuffer.toByteArray(), partNumber++); // 上传当前分块
chunkBuffer.reset(); // 重置缓冲区
}
}
// 上传剩余数据
if (chunkBuffer.size() > 0) {
uploadChunk(chunkBuffer.toByteArray(), partNumber);
}
```
2. **错误处理**:
- 捕获并处理管道流异常(如 `IOException`):
```java
try {
while ((bytesRead = pin.read(buffer)) != -1) { ... }
} catch (IOException e) {
// 处理读取中断(如写入端关闭)
s3AsyncClient.abortMultipartUpload(...); // 取消上传
}
```
### 注意事项
- **分块大小限制**:S3 要求除最后一块外,每块至少 5MB[^1]。
- **线程安全**:确保写入 `PipedOutputStream` 的线程和读取 `PipedInputStream` 的线程独立运行,避免死锁。
- **资源关闭**:上传完成后关闭所有流:
```java
pin.close();
pout.close();
```
### 替代方案:使用 `AsyncRequestBody.forBlockingInputStream()`
如果坚持使用单次上传(非分块),可尝试以下方法(SDK 2.40.0+):
```java
AsyncRequestBody requestBody = AsyncRequestBody.forBlockingInputStream(null); // 内容长度设为 null
CompletableFuture<?> writeFuture = CompletableFuture.runAsync(() -> {
try (OutputStream os = requestBody.outputStream()) {
// 将 PipedInputStream 数据复制到 os
byte[] buf = new byte[8192];
int bytesRead;
while ((bytesRead = pin.read(buf)) != -1) {
os.write(buf, 0, bytesRead);
}
} catch (IOException e) {
throw new UncheckedIOException(e);
}
}, executor);
s3AsyncClient.putObject(
PutObjectRequest.builder().bucket("bucket").key("key").build(),
requestBody
).thenRun(writeFuture::join); // 等待数据写入完成
```
> **警告**:此方法可能导致 SDK 内部使用大量内存,仅适用于小文件。
### 总结
| 方法 | 适用场景 | 风险 |
|--------------------|--------------------------|--------------------------|
| 分块上传 | 大文件/动态数据 | 实现复杂,需手动管理分块 |
| `forBlockingInputStream` | 小文件(<100MB) | 内存溢出风险 |
推荐优先使用 **分块上传方案** 处理未知大小的流数据[^1][^2]。
§§相关问题§§
1. 如何优化 AWS S3 分块上传的并发性能?
2. 使用 `PipedInputStream` 时如何避免死锁?
3. AWS SDK 中 `AsyncRequestBody` 的其他高级用法有哪些?
</think>
### AWS SDK 2.40.0 中从 PipedInputStream 上传未知大小数据到 S3 的解决方案
当无法预先计算数据流大小时,可以使用以下方法正确配置 `AsyncRequestBody` 实现安全上传:
#### 方法1:使用分块上传 (推荐)
```java
// 初始化分块上传
CreateMultipartUploadRequest createRequest = CreateMultipartUploadRequest.builder()
.bucket("your-bucket")
.key("object-key")
.build();
CreateMultipartUploadResponse response = s3AsyncClient.createMultipartUpload(createRequest).join();
String uploadId = response.uploadId();
// 分块读取并上传
List<CompletedPart> completedParts = new ArrayList<>();
byte[] buffer = new byte[8 * 1024 * 1024]; // 8MB块(与multipart配置匹配)
int partNumber = 1;
int bytesRead;
try (PipedInputStream pin = new PipedInputStream(pout, 1024*1024*10)) { // 10MB管道缓冲区
while ((bytesRead = pin.read(buffer)) != -1) {
// 创建分块请求体(明确提供当前块大小)
AsyncRequestBody partBody = AsyncRequestBody.fromBytes(Arrays.copyOf(buffer, bytesRead));
UploadPartRequest uploadRequest = UploadPartRequest.builder()
.bucket("your-bucket")
.key("object-key")
.uploadId(uploadId)
.partNumber(partNumber)
.build();
// 上传分块
UploadPartResponse partResponse = s3AsyncClient.uploadPart(uploadRequest, partBody).join();
completedParts.add(CompletedPart.builder()
.partNumber(partNumber)
.eTag(partResponse.eTag())
.build());
partNumber++;
}
}
// 完成上传
CompletedMultipartUpload completedUpload = CompletedMultipartUpload.builder()
.parts(completedParts)
.build();
s3AsyncClient.completeMultipartUpload(CompleteMultipartUploadRequest.builder()
.bucket("your-bucket")
.key("object-key")
.uploadId(uploadId)
.multipartUpload(completedUpload)
.build()).join();
```
#### 方法2:使用自定义AsyncRequestBody实现
```java
AsyncRequestBody requestBody = new AsyncRequestBody() {
@Override
public Optional<Long> contentLength() {
return Optional.empty(); // 明确表示长度未知
}
@Override
public void subscribe(Subscriber<? super ByteBuffer> subscriber) {
// 使用背压感知的发布者
new InputStreamConsumingPublisher(
new LengthCheckInputStream(pin), // 自定义包装流
Executors.newCachedThreadPool()
).subscribe(subscriber);
}
};
// 自定义流包装器避免LengthAwareInputStream检查
class LengthCheckInputStream extends InputStream {
private final InputStream wrapped;
public LengthCheckInputStream(InputStream in) {
this.wrapped = in;
}
@Override
public int read() throws IOException {
return wrapped.read();
}
@Override
public int read(byte[] b, int off, int len) throws IOException {
return wrapped.read(b, off, len);
}
@Override
public void close() throws IOException {
wrapped.close();
}
}
```
#### 关键配置要点:
1. **扩大管道缓冲区**:防止读写阻塞
```java
// 创建10MB缓冲区(默认仅1KB)
PipedInputStream pin = new PipedInputStream(pout, 10 * 1024 * 1024);
```
2. **使用独立线程写入**:避免死锁
```java
ExecutorService writerExecutor = Executors.newSingleThreadExecutor();
writerExecutor.submit(() -> {
try (PipedOutputStream pout = new PipedOutputStream(pin)) {
// 数据生成逻辑
pout.write(generateData());
}
});
```
3. **错误处理改进**:
```java
s3AsyncClient.putObject(putRequest, requestBody)
.whenComplete((response, exception) -> {
if (exception != null) {
// 重置管道或重启上传
pin = new PipedInputStream(1024*1024*10);
}
});
```
#### 性能优化建议:
1. **动态缓冲区调整**:
```java
int bufferSize = Math.max(8*1024*1024, pin.available());
byte[] dynamicBuffer = new byte[bufferSize];
```
2. **背压处理**:
```java
FlowControlUtil.addFlowControl(publisher, new FlowControlListener() {
@Override
public void onDrain() {
// 当接收端处理完毕时恢复写入
}
});
```
> **重要提示**:当使用 `PipedInputStream` 时,必须确保写入线程(`PipedOutputStream`)和读取线程同时运行,且写入线程需要在读取线程结束前保持活动状态。