【问题标题】:Multiple download requests with Alpakka S3 connector使用 Alpakka S3 连接器的多个下载请求
【发布时间】: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


【解决方案1】:

让我们通过剖析S3.download的返回类型来分解这里发生的事情:Source[Optional[JPair[Source[ByteString, NotUsed], ObjectMetadata]], NotUsed]

外部Source 代表未完成的请求。 如果在存储桶中找不到文件,Optional 为空。如果存在,它包含另一个 SourcePair 代表文件的字节内容,以及 ObjectMetadata 代表您正在下载的文件的元数据。

与直觉相反的是,Source 通常被表示为某个流式动作蓝图的一个冷酷、无状态、可共享的部分,只有在它具体化后才会焕发生机。对于外部Source,情况就是这样。然而,内心的Source却一反常态地瞬间“火爆”。一旦外部 Source 被具体化并发出一个项目,该项目代表一个打开的 HTTP 连接,您应该在(默认情况下)1 秒内开始使用,否则会引发 Response entity was not subscribed 错误。

在原始问题中,Source.combine 使用 Merge(_) 策略调用,这会导致并行实现。 Archive.zip 将按顺序处理文件,但如果完全消耗它收到的第一个 Source[ByteString] 需要超过 1 秒,第二个请求将在时间到来之前超时。

确保不会发生这种情况的一个可靠方法是消耗整个内部Source,然后再将其交给舞台中的下一个项目。考虑:

Source(1 to 10)
  .flatMapMerge(4, i => S3.download("my-s3-bucket", s"random$i.png")
    .log("started file download")
    .addAttributes(Attributes.logLevels(onElement = Attributes.LogLevels.Info))
    .flatMapConcat {
      case Some((s, m)) =>
        // for demo purposes, make sure individual downloads take >1 second
        s.delay(2.seconds, DelayOverflowStrategy.backpressure)
          // read entire file contents into a single ByteString
          .reduce(_ ++ _)
          .map(bs => (ArchiveMetadata(s"${UUID.randomUUID()}.png"), Source.single(bs)))
    })
  .log("completed file download")
  .addAttributes(Attributes.logLevels(onElement = Attributes.LogLevels.Info))
  .via(Archive.zip())
  .to(s3Sink)
  .run()

此(经过测试!)代码最多同时下载 4 个文件(flatMapMerge 的第一个参数)。请注意reduce 步骤在将响应传递给Archive.zip() 之前如何读取内存中的整个响应。这并不理想,但对于小文件来说可能是可以接受的。

【讨论】:

    猜你喜欢
    • 2020-09-08
    • 1970-01-01
    • 2019-09-19
    • 2020-05-19
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多