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

使用Alpakka 5.0.0压缩GCP存储桶文件夹时遇超时问题求助

问题

使用Alpakka 5.0.0开发Scala代码实现GCP存储桶内文件夹压缩功能时,第一个文件下载、压缩并上传完成后,处理第二个文件触发TimeoutException,提示响应实体未在1秒内订阅。观察到第二个请求似乎在第一个请求完全结束前就已启动,需求是实现文件的逐个流式处理。

原代码

import akka.actor.typed.ActorSystem
import akka.actor.typed.scaladsl.Behaviors
import akka.stream.Materializer
import akka.http.scaladsl.model.ContentType
import akka.http.scaladsl.model.MediaTypes.`application/zip`
import akka.stream.alpakka.file.ArchiveMetadata
import akka.stream.alpakka.file.scaladsl.Archive
import akka.stream.alpakka.googlecloud.storage.scaladsl.GCStorage
import akka.stream.scaladsl.Source

import java.util.UUID
import scala.concurrent.ExecutionContext

object GCPZipper extends App {

  implicit val system: ActorSystem[Nothing] = ActorSystem(Behaviors.empty, "zip-streamer")
  implicit val materializer: Materializer = Materializer(system)
  implicit val executionContext: ExecutionContext = system.executionContext

  private val bucket = "some-bucket-on-gcp"
  private val folder = "some-folder-on-bucket"
  private val contentType = ContentType(`application/zip`)

  val uuid = UUID.randomUUID().toString

  GCStorage.listBucket(bucket, Some(s"$folder/"))
    .log("zip-streamer - after list")
    .flatMapMerge(1, found => GCStorage.download(found.bucket, found.name).zip(Source.single(found)))
    .log("zip-streamer - after download")
    .collect({ case (Some(found), metadata) => (found, metadata) }) // to get rid of Option
    .log("zip-streamer - after collect")
    .map { case (found, metadata) => (ArchiveMetadata(metadata.name), found) }
    .log("zip-streamer - after archive metadata")
    .via(Archive.zip())
    .runWith(GCStorage.resumableUpload(bucket, s"${folder}_${uuid}_result.zip", contentType))    
}

异常信息

21:55:54.854 [zip-streamer-akka.actor.default-dispatcher-6] ERROR akka.stream.Materializer - [zip-streamer - after zip] Upstream failed.
java.util.concurrent.TimeoutException: Response entity was not subscribed after 1 second. Make sure to read the response `entity` body or call `entity.discardBytes()` on it -- in case you deal with `HttpResponse`, use the shortcut `response.discardEntityBytes()`. GET /storage/v1/b/some-bucket-on-gcp/o/643ffceb594787623f296844%2Fcog-hillshade-dtm-result.tiff Empty -> 200 OK Default(47874625 bytes)
    at akka.http.impl.engine.client.pool.SlotState$WaitingForResponseEntitySubscription.onTimeout(SlotState.scala:313)
    at akka.http.impl.engine.client.pool.NewHostConnectionPool$HostConnectionPoolStage$$anon$1$Event$.$anonfun$onTimeout$1(NewHostConnectionPool.scala:187)
    at akka.http.impl.engine.client.pool.NewHostConnectionPool$HostConnectionPoolStage$$anon$1$Event$.$anonfun$event0$1(NewHostConnectionPool.scala:189)
    at akka.http.impl.engine.client.pool.NewHostConnectionPool$HostConnectionPoolStage$$anon$1$Slot.runOneTransition$1(NewHostConnectionPool.scala:277)

Alpakka内部相关实现

// taken from DCStorageStram.scala

private def makeRequestSource[T: FromResponseUnmarshaller](request: Future[HttpRequest]): Source[T, NotUsed] =
  Source
    .fromMaterializer { (mat, attr) =>
      implicit val settings = resolveSettings(mat, attr)
      Source.lazyFuture { () =>
        request.flatMap { request =>
          GoogleHttp()(mat.system).singleAuthenticatedRequest[T](request)
        }(ExecutionContexts.parasitic)
      }
    }
    .mapMaterializedValue(_ => NotUsed)
问题根源

问题出在flatMapMerge(1, ...)的使用上:

  • 尽管设置了并行度为1,flatMapMerge仍会提前预取下一个元素对应的Source并发起订阅,导致第二个文件的HTTP请求已发送,但响应实体未被及时消费。
  • Akka HTTP客户端池默认对未订阅的响应实体设置1秒超时,超过时限后抛出TimeoutException。
  • Alpakka的GCStorage.download内部通过Source.lazyFuture发起请求,一旦被订阅就会触发HTTP请求,flatMapMerge的预订阅行为正好触发了这个提前请求的问题。
解决方案

将flatMapMerge替换为concatMap,concatMap会严格等待前一个Source完全处理完成(包括资源释放)后,才会订阅并启动下一个Source的处理,确保文件逐个串行处理,从根源避免提前请求导致的超时。

修改后的代码

import akka.actor.typed.ActorSystem
import akka.actor.typed.scaladsl.Behaviors
import akka.stream.Materializer
import akka.http.scaladsl.model.ContentType
import akka.http.scaladsl.model.MediaTypes.`application/zip`
import akka.stream.alpakka.file.ArchiveMetadata
import akka.stream.alpakka.file.scaladsl.Archive
import akka.stream.alpakka.googlecloud.storage.scaladsl.GCStorage
import akka.stream.scaladsl.Source

import java.util.UUID
import scala.concurrent.ExecutionContext

object GCPZipper extends App {

  implicit val system: ActorSystem[Nothing] = ActorSystem(Behaviors.empty, "zip-streamer")
  implicit val materializer: Materializer = Materializer(system)
  implicit val executionContext: ExecutionContext = system.executionContext

  private val bucket = "some-bucket-on-gcp"
  private val folder = "some-folder-on-bucket"
  private val contentType = ContentType(`application/zip`)

  val uuid = UUID.randomUUID().toString

  GCStorage.listBucket(bucket, Some(s"$folder/"))
    .log("zip-streamer - after list")
    // 替换flatMapMerge为concatMap,确保严格串行处理每个文件
    .concatMap(found => 
      GCStorage.download(found.bucket, found.name)
        .collect { case Some(source) => source } // 过滤掉None的情况
        .map(source => (ArchiveMetadata(found.name), source)) // 直接映射为压缩所需的元数据和流
    )
    .log("zip-streamer - after download and metadata mapping")
    .via(Archive.zip())
    .runWith(GCStorage.resumableUpload(bucket, s"${folder}_${uuid}_result.zip", contentType))    
}

修改说明

  1. 替换flatMapMerge为concatMap:彻底避免预订阅行为,保证前一个文件的下载、压缩流完全处理完成后,才会启动下一个文件的请求。
  2. 简化中间流操作:将原有的zip(Source.single(found))、collect、map步骤合并到concatMap内部,减少不必要的流转换,让逻辑更简洁高效。

内容的提问来源于stack exchange,提问作者Teliatko

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.20 03:15:04