【发布时间】:2021-03-17 12:18:13
【问题描述】:
所以,假设我有一个 Scala Vert.x Web REST API,它通过 HTTP 多部分请求接收文件上传。但是,它不会将传入的文件数据作为单个 InputStream 接收。相反,每个文件都是作为一系列字节缓冲区接收的,通过一些回调函数移交。
回调基本上是这样的:
// the callback that receives byte buffers (chunks) of the file being uploaded
// it is called multiple times until the full file has been received
upload.handler { buffer =>
// send chunk to backend
}
// the callback that gets called after the full file has been uploaded
// (i.e. after all chunks have been received)
upload.endHandler { _ =>
// do something after the file has been uploaded
}
// callback called if an exception is raised while receiving the file
upload.exceptionHandler { e =>
// do something to handle the exception
}
现在,我想使用这些回调将文件保存到 MinIO 存储桶中(如果您不熟悉,MinIO 基本上是自托管 S3,它的 API 与 S3 Java API 几乎相同) .
由于我没有文件句柄,我需要使用putObject() 将InputStream 放入MinIO。
我目前与 MinIO Java API 一起使用的低效解决方法如下所示:
// this is all inside the context of handling a HTTP request
val out = new PipedOutputStream()
val in = new PipedInputStream()
var size = 0
in.connect(out)
upload.handler { buffer =>
s.write(buffer.getBytes)
size += buffer.length()
}
upload.endHandler { _ =>
minioClient.putObject(
PutObjectArgs.builder()
.bucket("my-bucket")
.object("my-filename")
.stream(in, size, 50000000)
.build())
}
显然,这不是最佳选择。因为我在这里使用了一个简单的java.io 流,所以整个文件最终会被加载到内存中。
我不想在将文件放入对象存储之前将其保存到服务器上的磁盘。我想将它直接放入我的对象存储中。
如何使用 S3 API 和通过upload.handler 回调提供给我的一系列字节缓冲区来完成此操作?
编辑
我应该补充一点,我正在使用 MinIO,因为我不能使用商业托管的云解决方案,例如 S3。但是,正如 MinIO 网站上所提到的,我可以在使用 MinIO 作为我的存储解决方案的同时使用 Amazon 的 S3 Java SDK。
我尝试关注this guide on Amazon's website 将对象分块上传到 S3。
我尝试的解决方案如下所示:
context.request.uploadHandler { upload =>
println(s"Filename: ${upload.filename()}")
val partETags = new util.ArrayList[PartETag]
val initRequest = new InitiateMultipartUploadRequest("docs", "my-filekey")
val initResponse = s3Client.initiateMultipartUpload(initRequest)
upload.handler { buffer =>
println("uploading part", buffer.length())
try {
val request = new UploadPartRequest()
.withBucketName("docs")
.withKey("my-filekey")
.withPartSize(buffer.length())
.withUploadId(initResponse.getUploadId)
.withInputStream(new ByteArrayInputStream(buffer.getBytes()))
val uploadResult = s3Client.uploadPart(request)
partETags.add(uploadResult.getPartETag)
} catch {
case e: Exception => println("Exception raised: ", e)
}
}
// this gets called for EACH uploaded file sequentially
upload.endHandler { _ =>
// upload successful
println("done uploading")
try {
val compRequest = new CompleteMultipartUploadRequest("docs", "my-filekey", initResponse.getUploadId, partETags)
s3Client.completeMultipartUpload(compRequest)
} catch {
case e: Exception => println("Exception raised: ", e)
}
context.response.setStatusCode(200).end("Uploaded")
}
upload.exceptionHandler { e =>
// handle the exception
println("exception thrown", e)
}
}
}
这适用于小文件(我的测试小文件为 11 字节),但不适用于大文件。
对于大文件,upload.handler 中的进程会随着文件继续上传而逐渐变慢。此外,upload.endHandler 永远不会被调用,并且在文件上传 100% 后,文件会以某种方式继续上传。
但是,一旦我注释掉upload.handler 中的s3Client.uploadPart(request) 部分和upload.endHandler 中的s3Client.completeMultipartUpload 部分(基本上是丢弃文件而不是将其保存到对象存储中),文件上传就会进行正常并正确终止。
【问题讨论】:
-
我添加了尝试使用 AWS S3 Java API 将文件对象放入 MinIO。
-
zengularity.github.io/benji/s3/usage.html 正在使用 Akka-Stream 与任何 S3 兼容服务(AWS、Ceph、Minio)以流方式工作
-
好的。 没有 Akka使用它的任何例子?
-
是的,没错。我没有要求图书馆推荐。如果我无法从其他问题中获得任何收益,我为什么还要使用 Akka 来解决这个问题?
-
Akka, FS2 ... 问为什么不使用流媒体库进行流媒体传输对我来说至少很奇怪...
标签: java scala amazon-s3 vert.x minio