Scala函数式编程练习:如何实现全并行reduce方法?
并行实现Scala reduce方法的问题与解决
问题描述
我正在进行《Scala函数式编程》的练习,实现Scala中的reduce方法,但发现它并未并行运行。请问如何创建一个完全并行的实现?
我的reduce实现代码:
def reduce[A](pas: IndexedSeq[A])(f: (A, A) => A): Par[A] = if pas.isEmpty then throw new Exception("Can't reduce empty list") else if pas.size == 1 then unit(pas.head) else val (l, r) = pas.splitAt(pas.size / 2) reduce(l)(f).map2(reduce(r)(f))(f)
使用的Par类定义:
object Par: opaque type Par[A] = ExecutorService => Future[A] extension [A](pa: Par[A]) def run(s: ExecutorService): Future[A] = pa(s) def fork[A](a: => Par[A]): Par[A] = es => es.submit(new Callable[A] { def call = a(es).get }) def lazyUnit[A](a: => A): Par[A] = fork(unit(a))
问题分析与解决
你的实现没有并行执行的核心原因是:递归调用reduce(l)和reduce(r)时,两个任务都是在当前线程同步执行的,没有被提交到线程池进行异步调度。
要实现真正的并行,需要用fork方法将至少一个递归分支包装起来,强制其在独立线程中运行。修改后的代码如下:
def reduce[A](pas: IndexedSeq[A])(f: (A, A) => A): Par[A] = if pas.isEmpty then throw new Exception("Can't reduce empty list") else if pas.size == 1 then unit(pas.head) else val (l, r) = pas.splitAt(pas.size / 2) // 用fork包装左分支,触发并行执行 fork(reduce(l)(f)).map2(reduce(r)(f))(f) // 也可以同时包装左右两个分支:fork(reduce(l)(f)).map2(fork(reduce(r)(f)))(f)
关键说明
fork的作用是将任务提交到ExecutorService线程池异步执行,避免在当前线程同步阻塞计算。- 只需要包装其中一个分支,就能让左右两个递归任务并行处理——
map2会等待两个Par任务的结果,包装任意一个分支都会触发线程池的并行调度。同时包装两个分支也完全可行,不会产生副作用。 - 确保你的
map2实现是正确的:它需要能等待两个Future完成后再应用合并函数f,这是并行组合的基础(书中提供的标准map2实现已经处理了这一点)。
内容的提问来源于stack exchange,提问作者jorexe
相关产品推荐
相关产品推荐

