【问题标题】:Parallelize code in nested loops在嵌套循环中并行化代码
【发布时间】:2009-01-05 03:52:55
【问题描述】:

你总是听说函数式代码本质上比非函数式代码更容易并行化,所以我决定编写一个函数来执行以下操作:

给定一个字符串输入,合计每个字符串的唯一字符数。所以,给定输入[ "aaaaa"; "bbb"; "ccccccc"; "abbbc" ],我们的方法将返回a: 6; b: 6; c: 8

这是我写的:

(* seq<#seq<char>> -> Map<char,int> *)
let wordFrequency input =
    input
    |> Seq.fold (fun acc text ->
        (* This inner loop can be processed on its own thread *)
        text
        |> Seq.choose (fun char -> if Char.IsLetter char then Some(char) else None)
        |> Seq.fold (fun (acc : Map<_,_>) item ->
            match acc.TryFind(item) with
            | Some(count) -> acc.Add(item, count + 1)
            | None -> acc.Add(item, 1))
            acc
        ) Map.empty

这段代码在理想情况下是可并行化的,因为input 中的每个字符串都可以在其自己的线程上处理。它并不像看起来那么简单,因为内部循环将项目添加到所有输入之间共享的 Map。

我希望将内部循环分解到它自己的线程中,并且我不想使用任何可变状态。 如何使用异步工作流重写此函数?

【问题讨论】:

    标签: multithreading f# async-workflow


    【解决方案1】:

    你可以这样写:

    let wordFrequency =
      Seq.concat >> Seq.filter System.Char.IsLetter >> Seq.countBy id >> Map.ofSeq
    

    并仅使用两个额外字符将其并行化,以使用来自FSharp.PowerPack.Parallel.Seq DLL 的PSeq 模块,而不是普通的Seq 模块:

    let wordFrequency =
      Seq.concat >> PSeq.filter System.Char.IsLetter >> PSeq.countBy id >> Map.ofSeq
    

    例如,从 5.5Mb King James 圣经中计算频率所需的时间从 4.75s 下降到 0.66s。这在这台 8 核机器上是 7.2 倍的加速。

    【讨论】:

      【解决方案2】:

      正如已经指出的那样,如果您尝试让不同的线程处理不同的输入字符串,则会出现更新争用,因为每个线程都可以增加每个字母的计数。您可以让每个线程生成自己的地图,然后“将所有地图相加”,但最后一步可能会很昂贵(并且由于共享数据而不太适合使用线程)。我认为大型输入可能会使用如下算法运行得更快,其中每个线程处理不同的字母计数(对于输入中的所有字符串)。因此,每个线程都有自己独立的计数器,因此没有更新争用,也没有最后一步来组合结果。然而,我们需要预处理来发现“唯一字母集”,并且这一步确实存在相同的争用问题。 (在实践中,您可能预先知道字符的世界,例如字母,然后可以创建 26 个线程来处理 a-z,并绕过这个问题。)无论如何,大概问题主要是关于探索“如何编写 F#用于跨线程划分工作的异步代码,因此下面的代码演示了它。

      #light
      
      let input = [| "aaaaa"; "bbb"; "ccccccc"; "abbbc" |]
      
      // first discover all unique letters used
      let Letters str = 
          str |> Seq.fold (fun set c -> Set.add c set) Set.empty 
      let allLetters = 
          input |> Array.map (fun str -> 
              async { return Letters str })
          |> Async.Parallel 
          |> Async.Run     
          |> Set.union_all // note, this step is single-threaded, 
              // if input has many strings, can improve this
      
      // Now count each letter on a separate thread
      let CountLetter letter =
          let mutable count = 0
          for str in input do
              for c in str do
                  if letter = c then
                      count <- count + 1
          letter, count
      let result = 
          allLetters |> Seq.map (fun c ->
              async { return CountLetter c })
          |> Async.Parallel 
          |> Async.Run
      
      // print results
      for letter,count in result do
          printfn "%c : %d" letter count
      

      我确实“完全改变了算法”,主要是因为我的原始算法由于更新争用而不太适合直接数据并行化。根据您要学习的具体内容,这个答案可能会或可能不会让您特别满意。

      【讨论】:

      • Brian,这与我概述的算法相同吗?看起来好像。
      • 我不确定 Letters 方法是否比我对大型输入(每个包含几个 1000 个字符的字符串)采用的串行方法更有效。既然你问了,我并不是真的去学习具体的,只是写探索性代码:)
      • 最后,你是对的:原始代码的编写方式本质上是不可并行的,而将代码分开以进行适当并行化的唯一方法是以类似的方式重写您已经完成了上述操作。
      【解决方案3】:

      并行不等于异步,如Don Syme explains

      所以 IMO 你最好使用 PLINQ 进行并行化。

      【讨论】:

        【解决方案4】:

        我根本不会说 F#,但我可以解决这个问题。考虑使用 map/reduce:

        n = card(Σ) 为字母 Σ 中符号 σ 的数量。

        地图阶段:

        产生n个进程,其中第i个进程的分配是统计符号σi的出现次数> 在整个输入向量中。

        减少阶段

        按顺序收集每个 n 个进程的总数。那个向量就是你的结果。

        现在,此版本不会对串行版本产生任何改进;我怀疑这里有一个隐藏的依赖关系,这使得这本来就很难并行化,但我太累了,今晚无法证明这一点。

        【讨论】:

        • 有一点奇怪的依赖关系,对于非 F#'er 来说可能并不明显:Map 基本上是一个不可变的字典。当您向 Map 添加项目时,它会创建一个包含该项目的全新 Map 实例。换句话说,一个线程中的更改完全是(...)
        • (...) 与所有其他线程隔离。每个线程都可以返回自己的 Map 对象,然后我可以在最后总结结果。我不确定它是否会比非并行版本有更好的性能,我什至不确定如何在 F# 中实现它。
        猜你喜欢
        • 2021-12-22
        • 1970-01-01
        • 2021-08-13
        • 1970-01-01
        • 2020-09-27
        • 1970-01-01
        • 2017-06-20
        • 1970-01-01
        • 1970-01-01
        相关资源
        最近更新 更多