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

Scala中如何控制Future.sequence的并发执行数?

控制Scala Future并发数的几种实用方案

哈哈,这个坑我之前也踩过——直接用Future.sequence怼700个请求,差点把内部API服务器打挂😂。别担心,有好几种靠谱的方式能帮你把并发数控制在10个以内,给你挨个捋一遍:

1. 用Semaphore手动控制并发(无额外依赖)

这个方法利用Java的Semaphore做并发许可控制,同一时间最多允许10个任务获取许可执行,完成后释放许可给下一个任务。划重点:一定要把Future的创建逻辑包装成懒加载任务,不然所有Future会提前启动,Semaphore就白加了!

import scala.concurrent.{ExecutionContext, Future}
import java.util.concurrent.Semaphore

def processWithSemaphore[T](tasks: Seq[() => Future[T]], maxConcurrent: Int)(implicit ec: ExecutionContext): Future[Seq[T]] = {
  val semaphore = new Semaphore(maxConcurrent)
  
  val controlledFutures = tasks.map { task =>
    // 先获取许可,再执行任务,任务完成后无论成功失败都释放许可
    Future { semaphore.acquire() }
      .flatMap(_ => task())
      .andThen { case _ => semaphore.release() }
  }
  
  Future.sequence(controlledFutures)
}

// 用法示例:把你的每个REST请求包装成懒加载任务
// 假设seqOfFutures是你原本的Seq[Future[T]],要改成() => Future[T]避免提前执行
val lazyTasks = seqOfFutures.map(f => () => f)
val controlledResult = processWithSemaphore(lazyTasks, 10)
controlledResult.map(seqT => { /* 这里处理你的结果序列 */ })

2. 用Akka Streams(推荐,适合已有Akka的项目)

如果你的项目已经引入了Akka依赖,那Akka Streams的mapAsync算子简直是为这个场景量身定做的——它天生支持背压,能精准控制并发数,代码还特别简洁:

import akka.actor.ActorSystem
import akka.stream.scaladsl.{Sink, Source}
import scala.concurrent.Future

// 初始化ActorSystem(项目里如果已经有就不用重复创建了)
implicit val system: ActorSystem = ActorSystem("ConcurrentRequestControl")
import system.dispatcher

def processWithAkkaStreams[T](futures: Seq[Future[T]], maxConcurrent: Int): Future[Seq[T]] = {
  Source(futures)
    .mapAsync(maxConcurrent)(identity) // mapAsync的参数就是最大并发数
    .runWith(Sink.seq) // 收集所有结果成Seq[T]
}

// 用法
val controlledResult = processWithAkkaStreams(seqOfFutures, 10)
controlledResult.map(seqT => { /* 处理结果 */ })

这个方法的优势是:当某个任务完成后会立刻启动下一个,资源利用率更高,而且结果顺序和原序列一致,不用自己手动拼接。

3. 分批次执行(最简单的 fallback 方案)

如果不想加任何额外逻辑或依赖,也可以把700个请求分成每10个一批,等一批全部执行完再启动下一批。缺点是资源利用率不如前两种,但胜在简单易懂:

import scala.concurrent.{ExecutionContext, Future}

def processInBatches[T](futures: Seq[Future[T]], batchSize: Int)(implicit ec: ExecutionContext): Future[Seq[T]] = {
  futures.grouped(batchSize)
    .foldLeft(Future.successful(Seq.empty[T])) { (accumulatedResult, currentBatch) =>
      accumulatedResult.flatMap { results =>
        Future.sequence(currentBatch).map(results ++ _)
      }
    }
}

// 用法
val controlledResult = processInBatches(seqOfFutures, 10)
controlledResult.map(seqT => { /* 处理结果 */ })

总结

  • 已有Akka依赖:优先用Akka Streams,代码简洁高效
  • 无额外依赖:用Semaphore的方法,资源利用率更高
  • 追求极简:用分批次执行,适合快速实现

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.05.25 07:07:42