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
相关产品推荐
相关产品推荐

