【发布时间】:2020-04-03 18:25:27
【问题描述】:
我正在尝试使用Alpakka S3 connector 执行以下操作:
- 从 AWS S3 下载大量文件
- 通过 Alpakka Zip Archive Flow 流式传输下载的文件
- 将 Zip 流上传回 S3 Sink
我使用的代码是这样的:
val s3Sink: Sink[ByteString, Future[MultipartUploadResult]] = S3.multipartUpload("my-s3-bucket", "archive.zip")
val sourceList = (1 to 10).map(i => S3.download("my-s3-bucket", s"random$i.png").map {
case Some((s, m)) => (ArchiveMetadata(s"${UUID.randomUUID()}.png"), s)
})
val source = Source.combine(sourceList.head, sourceList.tail.head, sourceList.tail.tail: _*)(Merge(_))
source
.via(Archive.zip())
.to(s3Sink)
.run()
但是,这会导致以下错误:
Response entity was not subscribed after 1 second. Make sure to read the response entity body or call `discardBytes()` on it.
我怀疑这是因为 S3 连接器 expects every download response to be consumed 在移动到下一个连接器之前使用的底层 Akka Http,但是如果不引入等待/延迟,我无法以合理的方式处理这个问题。
我尝试使用带有bufferSize = 1 的队列,但这也不起作用。
我对 Akka 和 Akka Streams 还很陌生。
【问题讨论】:
-
我们正在解决这个问题。当我们这样做时,我会提供一个明确的答案。这可能是由于元素的内部缓冲,以及有限的连接池大小。您在后台使用同一个池进行上传和下载,因此如果发生争用,下载可能会在您有能力再次上传之前开始。
-
最终,我最终使用了来自普通 AWS SDK GetObjectRequest 的 InputStream,直到我想出一种使用 Alpakka/Akka Http 的方法
-
@LászlóvandenHoek 那方面有什么进展吗?我也遇到了这个问题
-
@AdamSzmyd 我们最终不需要进行任何 S3 到 S3 流式传输,所以我没有确切的答案。我确实记得我最初关于内部缓冲的断言是正确的,因此在准备好使用内部缓冲(响应字节)。原始问题中的
Source.combine太急切了。flatMapConcat会比单个文件名更好。 -
@AdamSzmyd 我自己添加了一个答案。
标签: scala amazon-s3 akka akka-stream alpakka