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

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的实现逻辑共同导致:

  1. Akka默认Dispatcher的线程池特性
    Akka默认Dispatcher基于ForkJoinPool,线程池大小默认与CPU核心数绑定(比如4核对应4个线程)。嵌套flatMap时,所有任务都在该线程池中执行。

  2. mapAsyncParallel的批量提交逻辑
    mapAsyncParallel会一次性将序列中所有元素转换为throttled Future并提交到线程池,每个throttled的第一步是Future { sem.acquire() }——这会直接占用线程池线程。

  3. 死锁形成的具体过程

    • 假设线程池有4个线程,mapAsync(4)设置信号量许可数为4。
    • 进入flatMap逻辑后,5个throttled Future被同时提交,前4个的sem.acquire()占用全部4个线程并成功获取许可。
    • 这4个线程接下来要执行f(current)(即Future.successful(scope))的回调,但此时线程池已无空闲线程,无法处理该Future的完成逻辑。
    • 第5个throttled Future的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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.07.01 12:04:57