Scala使用.par并行集合时单线程抛异常如何终止全部并行线程?
根因说明
Scala 标准库原生并行集合的默认异常处理逻辑为:等待所有已提交的并行任务全部执行完成后,才会将捕获到的异常汇总包装后向上抛出。你提供的示例中i=5对应的任务会进入永久休眠,因此并行集合会一直等待该任务结束,永远不会进入异常抛出阶段,最终表现为整体阻塞。
可实现需求的方案(保留.par使用的前提下)
- 方案1:自定义中断标识 + 自定义可中断线程池
并行集合默认使用的通用ForkJoinPool对中断的支持较弱,你可以自定义绑定线程池,配合全局可见的中断标识实现异常触发后全量任务终止。示例代码如下:
import java.util.concurrent.ForkJoinPool import scala.collection.parallel.ForkJoinTaskSupport // 线程可见的中断标识 @volatile private var isTaskCancelled: Boolean = false val testList = List(1, 2, 3, 4, 5).par // 绑定自定义线程池 val customPool = new ForkJoinPool(Runtime.getRuntime.availableProcessors()) testList.tasksupport = new ForkJoinTaskSupport(customPool) try { testList.foreach { i => // 每次执行业务逻辑前先检查是否已触发取消 if (isTaskCancelled) return @foreach println(s"i = $i") if (i == 5) { println("Sleeping forever") // 对阻塞操作捕获中断异常,响应取消信号 try { java.lang.Thread.sleep(Long.MaxValue) } catch { case _: InterruptedException => return @foreach } } throw new IllegalArgumentException("foo") } } catch { case e: IllegalArgumentException => // 捕获到业务异常后触发全局取消 isTaskCancelled = true // 强制关闭线程池,中断所有正在运行的任务 customPool.shutdownNow() // 重新抛出异常满足测试断言要求 throw e }
使用该方案需要注意:你需要在业务逻辑的关键执行节点主动检查isTaskCancelled标识,避免长时间无检查点的任务无法及时响应取消信号。
- 方案2:结果包裹检测法
如果不想修改默认线程池,也可以将每个并行任务的执行结果用Try包裹,遍历过程中一旦检测到失败结果就设置取消标识,后续的任务直接跳过执行。该方案的局限性是无法主动中断已经启动的长时间阻塞任务,仅适合任务逻辑有主动检查点、无永久阻塞逻辑的场景。
内容的提问来源于stack exchange,提问作者samthebest。
相关产品推荐
相关产品推荐

