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

Scala Futures推测执行:如何优先处理先返回的Future并分流慢Future

用Scala标准库轻松实现优先取先完成的Future需求

嗨,你这个场景其实Scala标准库里就有现成的简洁方案,完全不用手动去检查Future是否就绪~我给你两种常用的实现方式,都很优雅:

方式一:用Future.firstCompletedOf处理主路径,单独处理每个Future的后续逻辑

这个方法是最直接的,它会返回传入的Future集合中第一个完成的那个,不管是成功还是失败。你可以用它来处理主代码路径,同时给两个Future分别添加回调来处理后续逻辑。

举个完整的代码例子:

import scala.concurrent.{Future, ExecutionContext}
import scala.util.{Success, Failure}

// 模拟两个外部端点调用
val variableSpeedEndpoint: Future[String] = Future {
  // 模拟随机延迟,可能快可能慢
  Thread.sleep(scala.util.Random.nextInt(2000))
  "来自「速度可变型」的数据"
}

val reliableOldEndpoint: Future[String] = Future {
  Thread.sleep(1500) // 模拟固定的中等延迟
  "来自「可靠旧型」的数据"
}

// 必须提供隐式的ExecutionContext
implicit val ec: ExecutionContext = ExecutionContext.global

// 优先处理第一个完成的结果(主代码路径)
val firstDone = Future.firstCompletedOf(Seq(variableSpeedEndpoint, reliableOldEndpoint))
firstDone.onComplete {
  case Success(result) =>
    println(s"✅ 优先使用结果: $result")
    // 这里写你的主业务逻辑
  case Failure(ex) =>
    println(s"❌ 先返回的请求失败了: ${ex.getMessage}")
    // 主路径失败的 fallback 逻辑
}

// 同时处理两个端点的后续结果(不管是否被优先使用)
variableSpeedEndpoint.onComplete {
  case Success(data) => println(s"🔄 「速度可变型」返回结果,做后续处理: $data")
  case Failure(ex) => println(s"❌ 「速度可变型」请求失败: ${ex.getMessage}")
}

reliableOldEndpoint.onComplete {
  case Success(data) => println(s"🔄 「可靠旧型」返回结果,做后续处理: $data")
  case Failure(ex) => println(s"❌ 「可靠旧型」请求失败: ${ex.getMessage}")
}

方式二:用Future.select同时获取第一个完成的和剩余的Future

如果你想更精准地控制「第一个完成的」和「剩下的那个」的逻辑,Future.select会更合适——它返回的是一个Future[(Try[T], Iterable[Future[T]])],其中第一个元素是第一个完成的结果,第二个元素是剩下的未完成的Future集合。

示例代码:

import scala.concurrent.{Future, ExecutionContext}
import scala.util.{Success, Failure}

// 同样先定义两个端点Future
val variableSpeedEndpoint: Future[String] = Future {
  Thread.sleep(scala.util.Random.nextInt(2000))
  "来自「速度可变型」的数据"
}

val reliableOldEndpoint: Future[String] = Future {
  Thread.sleep(1500)
  "来自「可靠旧型」的数据"
}

implicit val ec: ExecutionContext = ExecutionContext.global

val selectedResult = Future.select(Seq(variableSpeedEndpoint, reliableOldEndpoint))
selectedResult.onComplete {
  case Success((firstResult, remainingFutures)) =>
    firstResult match {
      case Success(mainData) =>
        println(s"✅ 主路径使用结果: $mainData")
        // 处理剩下的那个Future
        remainingFutures.head.onComplete {
          case Success(remainingData) =>
            println(s"🔄 后续处理剩余结果: $remainingData")
          case Failure(ex) =>
            println(s"❌ 剩余请求失败: ${ex.getMessage}")
        }
      case Failure(ex) =>
        println(s"❌ 第一个完成的请求失败,尝试剩余请求")
        // 第一个失败了,用剩下的那个当备选
        remainingFutures.head.onComplete {
          case Success(fallbackData) =>
            println(s"✅ 备选路径使用结果: $fallbackData")
          case Failure(fallbackEx) =>
            println(s"❌ 所有请求都失败了: ${fallbackEx.getMessage}")
        }
    }
  case Failure(ex) =>
    println(s"❌ 意外错误: ${ex.getMessage}")
}

这两种方式都比手动检查Future状态简洁太多,完全利用Scala标准库的能力就能搞定你的需求~

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.19 09:55:39