Scala中Future列表在map内无法并行执行的原因排查
问题
我无法理解为何以下代码无法并行执行:
implicit val ec: ExecutionContext = system.dispatchers.lookup("db-non-blocking") println(s"Cores ${Runtime.getRuntime().availableProcessors()}") def print(i: Int): Future[Unit] = Future { if (i == 1) { Thread.sleep(5000) } println(s"Finished $i at ${new Date(System.currentTimeMillis())} (Thread ${Thread.currentThread().getName})") } Future { println(s"${new Date(System.currentTimeMillis())} (Thread ${Thread.currentThread().getName})") Seq(1, 2, 3, 4, 5) }.map(d => { d.map(u => print(u)) })
运行代码后输出如下:
Cores 16 Tue Jun 18 14:19:15 CEST 2024 (Thread plugin-system-db-non-blocking-16) Finished 1 at Tue Jun 18 14:19:20 CEST 2024 (Thread plugin-system-db-non-blocking-16) Finished 2 at Tue Jun 18 14:19:20 CEST 2024 (Thread plugin-system-db-non-blocking-16) Finished 3 at Tue Jun 18 14:19:20 CEST 2024 (Thread plugin-system-db-non-blocking-16) Finished 4 at Tue Jun 18 14:19:20 CEST 2024 (Thread plugin-system-db-non-blocking-16) Finished 5 at Tue Jun 18 14:19:20 CEST 2024 (Thread plugin-system-db-non-blocking-16)
第一个Future完成后,map返回一个Future列表,但这些创建的Future并未异步执行。所有线程都使用同一个Dispatcher,该调度器的线程池大小(固定50)足以支持5个线程并行,但所有Future似乎都在同一个线程运行。调度器配置如下:
dispatcher-non-blocking { type = Dispatcher executor = "thread-pool-executor" thread-pool-executor { fixed-pool-size = 50 } throughput = 1 }
为何这些Future无法并行且都在ID为16的线程运行?更换调度器则可并行,若不将Seq包裹在Future内也可并行,为何在map内会同步执行?
解答
核心原因:map回调的执行特性
Scala/Akka中Future.map的回调逻辑,是在完成原Future的线程上同步执行的,除非回调内部显式切换执行上下文。
具体到你的代码:
- 第一个
Future执行完毕后,它的map回调会直接在执行该Future的线程(即plugin-system-db-non-blocking-16)上启动。 - 你在
map里调用d.map(u => print(u)),这是普通的集合遍历操作——整个遍历过程是同步的,会在当前线程(16号线程)上依次调用print(u)。 print(u)创建的Future虽然会提交到调度器,但由于当前线程正处于map回调的执行流程中,调度器会优先让当前线程处理这些新提交的任务,最终导致所有print里的逻辑都在同一个线程串行执行。
为什么其他场景能并行?
- 不包裹在Future的map里时:直接在主线程(或非Future回调线程)调用
Seq(1,2,3,4,5).map(print),每个print创建的Future会被立即提交到调度器,线程池可以分配不同空闲线程并行处理任务。 - 更换调度器后能并行:不同调度器的线程池策略、任务调度逻辑存在差异,部分调度器会强制将新任务分配到空闲线程,而非复用当前执行回调的线程。
解决方法
要让这些Future并行执行,需避免集合遍历和Future创建的过程阻塞在原Future的回调线程上,常见方案有两种:
- 使用
Future.traverse替代手动遍历:它会异步遍历集合并并行执行每个元素对应的Future任务,自动利用执行上下文的线程池:Future { println(s"${new Date(System.currentTimeMillis())} (Thread ${Thread.currentThread().getName})") Seq(1, 2, 3, 4, 5) }.flatMap(d => Future.traverse(d)(print)) - 显式将遍历逻辑切换到执行上下文:通过
Future包裹遍历过程,让遍历和Future创建在新线程上执行:Future { println(s"${new Date(System.currentTimeMillis())} (Thread ${Thread.currentThread().getName})") Seq(1, 2, 3, 4, 5) }.flatMap(d => Future { d.map(print) })
内容的提问来源于stack exchange,提问作者xSmorpheusSx
相关产品推荐
相关产品推荐

