【问题标题】:Attemping to parallelize a nested loop in Scala尝试在 Scala 中并行化嵌套循环
【发布时间】:2019-05-27 07:44:05
【问题描述】:

我正在使用嵌套循环和外部 jar 比较 scala/spark 中的 2 个数据帧。

for (nrow <- dfm.rdd.collect) {   
  var mid = nrow.mkString(",").split(",")(0)
  var mfname = nrow.mkString(",").split(",")(1)
  var mlname = nrow.mkString(",").split(",")(2)  
  var mlssn = nrow.mkString(",").split(",")(3)  

  for (drow <- dfn.rdd.collect) {
    var nid = drow.mkString(",").split(",")(0)
    var nfname = drow.mkString(",").split(",")(1)
    var nlname = drow.mkString(",").split(",")(2)  
    var nlssn = drow.mkString(",").split(",")(3)  

    val fNameArray = Array(mfname,nfname)
    val lNameArray = Array (mlname,nlname)
    val ssnArray = Array (mlssn,nlssn)

    val fnamescore = Main.resultSet(fNameArray)
    val lnamescore = Main.resultSet(lNameArray)
    val ssnscore =  Main.resultSet(ssnArray)

    val overallscore = (fnamescore +lnamescore +ssnscore) /3

    if(overallscore >= .95) {
       println("MeditechID:".concat(mid)
         .concat(" MeditechFname:").concat(mfname)
         .concat(" MeditechLname:").concat(mlname)
         .concat(" MeditechSSN:").concat(mlssn)
         .concat(" NextGenID:").concat(nid)
         .concat(" NextGenFname:").concat(nfname)
         .concat(" NextGenLname:").concat(nlname)
         .concat(" NextGenSSN:").concat(nlssn)
         .concat(" FnameScore:").concat(fnamescore.toString)
         .concat(" LNameScore:").concat(lnamescore.toString)
         .concat(" SSNScore:").concat(ssnscore.toString)
         .concat(" OverallScore:").concat(overallscore.toString))
    }
  }
}

我希望做的是为外循环添加一些并行性,这样我就可以创建一个 5 个线程池并从外循环的集合中提取 5 条记录,并将它们与内循环的集合进行比较,而不是连续执行此操作。所以结果是我可以指定线程数,在任何给定时间从外部循环的集合处理中针对内部循环中的集合处理 5 条记录。我该怎么做呢?

【问题讨论】:

    标签: scala apache-spark dataframe parallel-processing parallel-collections


    【解决方案1】:

    让我们从分析您在做什么开始。您将dfm 的数据收集给驱动程序。然后,对于您从 dfn 收集数据的每个元素,对其进行转换并计算每对元素的分数。

    这在很多方面都存在问题。首先,即使不考虑并行计算,对dfn 的元素的转换也与dfm 的元素一样多次。此外,您为dfm 的每一行收集dfn 的数据。这是很多网络通信(驱动程序和执行程序之间)。

    如果您想使用 spark 并行计算,您需要使用 API(RDD、SQL 或数据集)。您似乎想使用 RDD 来执行笛卡尔积(这是 O(N*M),所以要小心,这可能需要一段时间)。

    让我们首先在笛卡尔积之前转换数据,以避免每个元素多次执行它们。另外,为了清楚起见,让我们定义一个案例类来包含您的数据和一个将您的数据帧转换为该案例类的 RDD 的函数。

    case class X(id : String, fname : String, lname : String, lssn : String)
    def toRDDofX(df : DataFrame) = {
        df.rdd.map(row => {
            // using pattern matching to convert the array to the case class X
            row.mkString(",").split(",") match {
                case Array(a, b, c, d) => X(a, b, c, d)
            } 
        })
    }
    

    然后,我使用filter 仅保留分数超过.95 的元组,但您可以使用mapforeach...,具体取决于您打算做什么。

    val rddn = toRDDofX(dfn)
    val rddm = toRDDofX(dfm)
    rddn.cartesian(rddm).filter{ case (xn, xm) => {
        val fNameArray = Array(xm.fname,xn.fname)
        val lNameArray = Array(xm.lname,xn.lname)
        val ssnArray = Array(xm.lssn,xn.lssn)
    
        val fnamescore = Main.resultSet(fNameArray)
        val lnamescore = Main.resultSet(lNameArray)
        val ssnscore =  Main.resultSet(ssnArray)
    
        val overallscore = (fnamescore +lnamescore +ssnscore) /3
        // and then, let's say we filter by score
        overallscore > .95
    }} 
    

    【讨论】:

    • 非常感谢这个详细的解释。您的代码正在运行,通过将其设置为 val rdd 然后 rdd.take(100).foreach(println) 它将过滤后的记录显示为: X(val1, val2,.....) 我现在正试图弄清楚如何将此数组 RDD 写入数据帧,以便我可以将其写回数据库。不断出错,但我一直在玩它。再次TY!
    • 我想出了如何将结果写回数据框。谢谢!
    【解决方案2】:

    这不是迭代 spark 数据帧的正确方法。主要关注的是dfm.rdd.collect。如果数据框任意大,您最终会出现异常。这是因为collect 函数本质上将所有数据都带入了主节点。

    另一种方法是使用 rdd 的 foreach 或 map 构造。

    dfm.rdd.foreach(x => {
        // your logic
    }  
    

    现在您正尝试在此处迭代第二个数据帧。恐怕那是不可能的。优雅的方法是加入 dfmdfn 并迭代生成的数据集以计算您的函数。

    【讨论】:

    • Scala/spark 的新手,所以我不知道所有深奥的来龙去脉。这就是我在这里问这个问题的原因......因此不明白需要否决票。你会如何建议我实现数据框连接代码来完成这个?
    • 首先,很抱歉投反对票,我没有这样做。在继续解决问题之前,您可能需要更深入地研究火花。我将继续的方式是将 dfm 和 dfn 加入它们各自的 uniqueId 上。然后迭代生成的数据帧以生成另一个具有所需字段的数据帧。 stackoverflow.com/questions/49252670/… 可能会有所帮助。另一个链接jaceklaskowski.gitbooks.io/mastering-spark-sql/…
    • 我会查看您建议的链接。不幸的是,我无法加入 ID,我必须在两个数据帧上执行笛卡尔运算,以将外部循环中的每条记录与内部循环中的每条记录进行比较。
    • 如果你的数据集很小(GB),你可以尝试做笛卡尔。看到问题是你不能在火花中对数据帧进行嵌套循环。方法是join和foreach/map。
    猜你喜欢
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 2020-09-27
    • 1970-01-01
    • 2017-06-20
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多