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

