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

Scala带超时等待Future序列,超时不抛TimeoutException且未完成视为失败

实现方案

你之前的代码出现TimeoutException的核心原因是:Future.sequence需要等所有输入的Future全部执行完成才会生成最终的总Future,你直接对这个总Future调用Await.result设置超时,只要有任意一个Future没有在超时时间内跑完,就会触发Await的超时异常,和你提前把Future封装成Try没有关系。

要实现到点立刻返回、未完成的Future统一视为失败的需求,核心思路是给每个独立的Future绑定超时逻辑,所有Future都会在你指定的超时时间内返回结果(正常执行结果/超时失败),这样总Future自然会在超时时间内完成,不会触发Await的超时。

实现代码

import scala.concurrent.{Future, ExecutionContext, Promise}
import scala.concurrent.duration.FiniteDuration
import scala.util.{Success, Failure, Try}
import java.util.concurrent.{Executors, TimeUnit}

// 全局单例定时调度器,避免反复创建线程
private val timeoutScheduler = Executors.newSingleThreadScheduledExecutor()

/**
 * 给单个Future添加超时逻辑,超时未完成则返回超时失败的Try
 */
private def withTimeout[T](fut: Future[T], timeout: FiniteDuration)(implicit ex: ExecutionContext): Future[Try[T]] = {
  val resPromise = Promise[Try[T]]()
  // 注册超时任务:到点直接返回超时失败
  val timeoutTask = new Runnable {
    override def run(): Unit = {
      resPromise.trySuccess(Failure(new java.util.concurrent.TimeoutException("Future执行超时")))
    }
  }
  val scheduled = timeoutScheduler.schedule(timeoutTask, timeout.toMillis, TimeUnit.MILLISECONDS)

  // 原Future执行完成后处理结果,取消超时任务
  fut.onComplete { tryRes =>
    scheduled.cancel(false)
    resPromise.trySuccess(tryRes)
  }

  resPromise.future
}

/**
 * 等待一组Future执行,到超时时间后立刻返回所有结果,未完成的Future视为超时失败
 * @param futures 待执行的Future序列
 * @param timeout 最长等待时长
 * @return 序列中每个Future的执行结果,Success为执行成功,Failure为执行失败/超时
 */
def waitAllWithTimeout[T](futures: Seq[Future[T]], timeout: FiniteDuration)(implicit ex: ExecutionContext): Future[Seq[Try[T]]] = {
  Future.sequence(futures.map(f => withTimeout(f, timeout)))
}

调用示例

import scala.concurrent.duration._
import scala.concurrent.ExecutionContext.Implicits.global

val futures: Seq[Future[MyObject]] = ...
val waitTimeout = waitTimeoutMillis.millis
// 这里Await的超时可以比实际业务超时稍长一点,避免调度误差导致意外超时
val segments = Await.result(waitAllWithTimeout(futures, waitTimeout), waitTimeout + 100.millis)

说明

  • 返回的Seq[Try[T]]中,Success对应执行成功的Future,Failure如果是TimeoutException类型表示该Future超时未完成,其他异常类型表示该Future执行过程中抛出错误
  • 如果你是短期运行的程序,程序退出前记得调用timeoutScheduler.shutdown()关闭定时调度线程池,避免JVM无法正常退出;如果是长期运行的服务可以保留该调度器复用。

内容的提问来源于stack exchange,提问作者Or Bar Yaacov

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.09.25 08:15:04