使用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)) }
修改说明
- 替换
flatMapMerge为concatMap:彻底避免预订阅行为,保证前一个文件的下载、压缩流完全处理完成后,才会启动下一个文件的请求。 - 简化中间流操作:将原有的
zip(Source.single(found))、collect、map步骤合并到concatMap内部,减少不必要的流转换,让逻辑更简洁高效。
内容的提问来源于stack exchange,提问作者Teliatko
相关产品推荐
相关产品推荐

