Scala中mapAsync独立可用但嵌套flatMap时阻塞问题排查
问题:mapAsync嵌套flatMap时阻塞
实现了一个通用序列的mapAsync方法,支持限制并发数,独立调用时运行正常,但将其嵌套在Future.flatMap中调用时程序会卡住。
示例代码
import akka.actor.ActorSystem import java.util.concurrent.Semaphore import scala.collection.generic.CanBuildFrom import scala.concurrent.duration.Duration import scala.concurrent.{Await, ExecutionContext, Future} object Test extends App { implicit class SeqWrapper[+A, M[X] <: TraversableOnce[X]](underlying: M[A]) { def mapAsyncSequential[B]( f: A ⇒ Future[B] )(implicit ec: ExecutionContext, cbf: CanBuildFrom[M[Future[B]], B, M[B]]): Future[M[B]] = { underlying .foldLeft(Future.successful(cbf())) { (fr, fa) ⇒ for { r ← fr a ← f(fa) } yield r += a } .map(_.result()) } def mapAsyncParallel[B]( f: A ⇒ Future[B] )(implicit ec: ExecutionContext, cbf: CanBuildFrom[M[Future[B]], B, M[B]]): Future[M[B]] = { underlying .map(f) .foldLeft(Future.successful(cbf())) { (fr, fb) ⇒ for { r ← fr b ← fb } yield (r += b) } .map(_.result()) } def mapAsync[B]( concurrentInstances: Int )(f: A ⇒ Future[B])(implicit ec: ExecutionContext, cbf: CanBuildFrom[M[Future[B]], B, M[B]]): Future[M[B]] = { val sem = new Semaphore(concurrentInstances) def throttled(current: A): Future[B] = Future { println("acquiring") sem.acquire() println("acquired") }.flatMap { _ ⇒ println("operation") f(current).andThen { case _ ⇒ println("release") sem.release() } } mapAsyncParallel(throttled) } } implicit class FutureValueProvider[T](f: Future[T]) { def futureValue: T = Await.result(f, Duration.Inf) } implicit val system: ActorSystem = ActorSystem("test") import system.dispatcher try { def f: Future[List[Int]] = { List(1, 2, 3, 4, 5).mapAsync(4) { scope => Future.successful(scope) } } println(f.futureValue) // 正常运行 println(Future.successful(1).flatMap(_ => f).futureValue) // 程序卡住 } finally system.terminate() }
现象
- 独立调用时输出正常:
acquiring acquiring acquired acquiring acquired operation operation acquiring acquired acquiring acquired operation operation release release acquired release release operation release List(1, 2, 3, 4, 5)
- 嵌套flatMap调用时程序卡住,输出停留在:
acquiring acquired acquiring acquired acquiring acquired acquiring acquired acquiring
原因分析
核心问题是执行上下文线程池耗尽引发的死锁,由Semaphore的使用方式、mapAsyncParallel的实现逻辑共同导致:
Akka默认Dispatcher的线程池特性
Akka默认Dispatcher基于ForkJoinPool,线程池大小默认与CPU核心数绑定(比如4核对应4个线程)。嵌套flatMap时,所有任务都在该线程池中执行。mapAsyncParallel的批量提交逻辑mapAsyncParallel会一次性将序列中所有元素转换为throttledFuture并提交到线程池,每个throttled的第一步是Future { sem.acquire() }——这会直接占用线程池线程。死锁形成的具体过程
- 假设线程池有4个线程,
mapAsync(4)设置信号量许可数为4。 - 进入
flatMap逻辑后,5个throttledFuture被同时提交,前4个的sem.acquire()占用全部4个线程并成功获取许可。 - 这4个线程接下来要执行
f(current)(即Future.successful(scope))的回调,但此时线程池已无空闲线程,无法处理该Future的完成逻辑。 - 第5个
throttledFuture的sem.acquire()等待信号量许可,而前4个线程因无法完成f(current)的回调,永远不会调用sem.release()——信号量无法释放,线程池也无法腾出资源,最终形成死锁。
- 假设线程池有4个线程,
独立调用正常的原因
独立调用时,Await.result会触发ForkJoinPool的线程补偿机制:当线程池线程全部阻塞时,会临时新增线程处理任务,让f(current)的Future得以执行,最终触发sem.release()释放许可,保证流程继续。但嵌套flatMap时,逻辑完全在Future回调链中,线程补偿机制无法触发,导致线程池彻底耗尽后无法恢复。
修复思路
- 避免用线程池线程执行阻塞操作:
sem.acquire()是阻塞操作,不应放在Future计算体内,可改用非阻塞信号量实现,或用tryAcquire配合Future.recoverWith重试,避免线程被占用。 - 调整并发任务提交逻辑:不要一次性提交所有Future,而是按需提交当前可处理的任务,避免瞬间占满线程池。
- 使用成熟的限流工具:比如Akka Streams的
mapAsync算子,这类工具已内置处理并发控制和线程池调度的逻辑,无需手动实现。
内容的提问来源于stack exchange,提问作者prongs
相关产品推荐
相关产品推荐

