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

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

相关产品推荐
方舟 Agent Plan

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

最近更新时间:2026.06.21 20:52:09