使用Future.sequence失败后如何清理已创建的容器资源?
如何在Scala Future失败时清理已创建的Container资源
假设我们有一个在堆上分配资源的Container类,它的创建过程耗时,因此封装在Future中:
class Container(val n: Int) { allocateMemoryOnTheHeap(n) def cleanup(): Unit = ??? } def prepareContainer(n: Int): Future[Container] = ???
我们的任务是下载一系列值并将其存入该容器类中。若任一下载操作失败,需要返回失败结果并释放已创建容器的资源。
当前的实现代码如下:
def downloadNumber(): Future[Int] = ??? def numbersProvider(howMany: Int): Future[Seq[Int]] = Future.sequence((0 until howMany).map(_ => downloadNumber())) def prepareContainersForValuesFromUpstream(): Future[Seq[Container]] = { val futureNumbers = numbersProvider(42) val containersFuture: Future[Seq[Future[Container]]] = futureNumbers.map { numbers => numbers.map(prepareContainer(_)) } containersFuture.flatMap(Future.sequence) } def main() = { val containersFuture = prepareContainersForValuesFromUpstream() val result = Await.result(containersFuture, timeout) result match { case Success(_) => println("下载完成,值已存入容器以备后续使用!") case Failure(_) => // TODO: 如何释放已创建容器的资源? } }
当前实现的问题在于:使用Future.sequence批量处理时,一旦某个Future失败,sequence会直接返回失败结果,无法获取已经成功创建的Container实例,自然无法调用cleanup方法。如果忽略失败的Future只收集成功结果,又不符合“任一下载失败就返回失败”的需求。
解决方案:逐个处理并追踪已创建容器
核心思路是放弃直接用Future.sequence批量处理,改为逐个处理每个任务,同时维护已创建容器的列表,一旦任何步骤失败,立即遍历列表清理所有已分配的资源。
实现代码
import scala.concurrent.{Future, Await} import scala.concurrent.ExecutionContext.Implicits.global import scala.util.{Success, Failure, Try} import scala.concurrent.duration._ class Container(val n: Int) { // 模拟堆内存分配逻辑 private val allocatedData = new Array[Byte](n * 1024 * 1024) def cleanup(): Unit = { println(s"清理Container(${n})占用的资源") // 这里可添加实际资源释放逻辑,比如置空引用、调用底层释放API等 } } def prepareContainer(n: Int): Future[Container] = Future { // 模拟耗时的容器创建过程 Thread.sleep(100) new Container(n) } def downloadNumber(): Future[Int] = Future { // 模拟随机失败的下载操作 val num = scala.util.Random.nextInt(10) if (num == 5) throw new RuntimeException("下载失败") num } def numbersProvider(howMany: Int): Future[Seq[Int]] = Future.sequence((0 until howMany).map(_ => downloadNumber())) def prepareContainersSafely(numbers: Seq[Int]): Future[Seq[Container]] = { // 使用foldLeft逐步构建容器列表,同时追踪操作状态 numbers.foldLeft(Future.successful((Seq.empty[Container], true))) { case (accFuture, num) => accFuture.flatMap { case (containers, isSuccess) => if (!isSuccess) { // 之前步骤已失败,直接返回当前状态 Future.successful((containers, false)) } else { prepareContainer(num).map { container => (containers :+ container, true) }.recoverWith { case ex => // 失败时清理所有已创建的容器 containers.foreach(_.cleanup()) Future.failed(ex) } } } }.map(_._1) } def prepareContainersForValuesFromUpstream(): Future[Seq[Container]] = { numbersProvider(42).flatMap { numbers => prepareContainersSafely(numbers) }.recoverWith { case ex => // 下载数字阶段失败时,无容器需要清理,直接返回失败 Future.failed(ex) } } def main(): Unit = { val timeout = 10.seconds Try(Await.result(prepareContainersForValuesFromUpstream(), timeout)) match { case Success(_) => println("下载完成,值已存入容器以备后续使用!") case Failure(ex) => println(s"操作失败:${ex.getMessage}") } }
方案说明
prepareContainersSafely方法:通过foldLeft逐个处理每个数字,维护包含已创建容器和操作状态的元组。每一步创建容器成功则添加到列表,失败则立即清理所有已创建资源并返回失败。- 错误传播逻辑:一旦任何步骤失败,后续步骤不再执行,直接返回失败结果,同时保证已分配的资源被完全清理。
- 下载阶段失败处理:如果
numbersProvider本身失败(比如下载数字时出错),此时还未创建任何容器,直接返回失败即可。
更简洁的通用实现:自定义safeTraverse
如果觉得foldLeft的写法不够直观,可以实现一个通用的safeTraverse方法,复用在需要资源清理的场景中:
def safeTraverse[A, B <: { def cleanup(): Unit }](seq: Seq[A])(f: A => Future[B]): Future[Seq[B]] = { seq.foldLeft(Future.successful(Seq.empty[B])) { (acc, a) => acc.flatMap { bs => f(a).map(b => bs :+ b).recoverWith { ex => // 失败时清理所有已成功创建的实例 bs.foreach(_.cleanup()) Future.failed(ex) } } } } // 修改prepareContainersForValuesFromUpstream方法 def prepareContainersForValuesFromUpstream(): Future[Seq[Container]] = { numbersProvider(42).flatMap { numbers => safeTraverse(numbers)(prepareContainer) } }
这个safeTraverse方法可以处理任何实现了cleanup方法的资源类,通用性更强。
内容的提问来源于stack exchange,提问作者delabania
相关产品推荐
相关产品推荐

