【问题标题】:Joining process with broadcast variable ends up endless spilling加入广播变量的过程会导致无休止的溢出
【发布时间】:2015-10-05 22:10:56
【问题描述】:

我正在以独立模式从文本文件中加入两个 RDD。一个有 4 亿 (9 GB) 行,另一个有 400 万 (110 KB)。

3-grams  doc1           3-grams   doc2
ion -    100772C111      ion -    200772C222  
on  -    100772C111      gon -    200772C222  
 n  -    100772C111        n -    200772C222
... -    ....            ... -    .... 
ion -    3332145654      on  -    58898874
mju -    3332145654      mju -    58898874
... -    ....            ... -    ....

在每个文件中,文档编号(doc1 或 doc2)一个接一个地出现。作为加入的结果,我想在文档之间获得一些常见的 3-grams.e.g.

  (100772C111-200772C222,2) --> There two common 3-grams which are 'ion' and  '  n'

我运行代码的服务器有 128 GB RAM 和 24 个内核。我使用 -Xmx64G 设置了我的 IntelliJ 配置 - VM 选项部分

这是我的代码:

val conf = new SparkConf().setAppName("abdulhay").setMaster("local[4]").set("spark.shuffle.spill", "true")
      .set("spark.shuffle.memoryFraction", "0.6").set("spark.storage.memoryFraction", "0.4")
      .set("spark.executor.memory","40g")
      .set("spark.driver.memory","40g")

val sc = new SparkContext(conf)

val emp = sc.textFile("\\doc1.txt").map(line => (line.split("\t")(3),line.split("\t")(1))).distinct()
    val emp_new = sc.textFile("\\doc2.txt").map(line => (line.split("\t")(3),line.split("\t")(1))).distinct()

val emp_newBC = sc.broadcast(emp_new.groupByKey.collectAsMap)

val joined = emp.mapPartitions(iter => for {
      (k, v1) <- iter
      v2 <- emp_newBC.value.getOrElse(k, Iterable())
    } yield (s"$v1-$v2", 1))

val olsun = joined.reduceByKey((a,b) => a+b)

olsun.map(x => x._1 + "\t" + x._2).saveAsTextFile("...\\out.txt")

如所见,在使用广播变量的加入过程中,我的键值发生了变化。所以看来我需要重新分区连接的值?而且它非常昂贵。结果,我结束了太多溢出的问题,并且从未结束。我认为 128 GB 内存必须足够。据我了解,当使用广播变量时​​,洗牌正在显着减少?那么我的申请有什么问题呢?

提前致谢。

编辑:

我也试过spark的join功能如下:

var joinRDD = emp.join(emp_new);

val kkk = joinRDD.map(line => (line._2,1)).reduceByKey((a, b) => a + b)

又一次溢出太多。

EDIT2:

val conf = new SparkConf().setAppName("abdulhay").setMaster("local[12]").set("spark.shuffle.spill", "true")
      .set("spark.shuffle.memoryFraction", "0.4").set("spark.storage.memoryFraction", "0.6")
      .set("spark.executor.memory","50g")
      .set("spark.driver.memory","50g")
    val sc = new SparkContext(conf)

val emp = sc.textFile("S:\\Staff_files\\Mehmet\\Projects\\SPARK - scala\\wos14.txt").map{line => val s = line.split("\t"); (s(5),s(0))}//.distinct()
    val emp_new = sc.textFile("S:\\Staff_files\\Mehmet\\Projects\\SPARK - scala\\fwo_word.txt").map{line => val s = line.split("\t"); (s(3),s(1))}//.distinct()

    val cog = emp_new.cogroup(emp)

val skk =  cog.flatMap {
      case (key: String, (l1: Iterable[String], l2: Iterable[String])) =>
        (l1.toSeq ++ l2.toSeq).combinations(2).map { case Seq(x, y) => if (x < y) ((x, y),1) else ((y, x),1) }.toList
    }

    val com = skk.countByKey()

【问题讨论】:

    标签: scala join intellij-idea apache-spark


    【解决方案1】:

    我不会使用广播变量。当你说:

    val emp_newBC = sc.broadcast(emp_new.groupByKey.collectAsMap)
    

    Spark 首先将整个数据集移动到主节点中,这是一个巨大的瓶颈,并且容易在主节点上产生内存错误。然后,这块内存被洗牌回到所有节点(大量的网络开销),也必然会在那里产生内存问题。

    相反,使用 join 加入 RDD 本身(参见描述 here

    还要弄清楚你的键是否太少。加入 Spark 基本上需要将整个 key 加载到内存中,如果你的 key 太少,对于任何给定的 executor 来说,分区可能仍然太大。

    单独说明:reduceByKey 无论如何都会重新分区。

    编辑:---------

    好的,鉴于澄清,并假设每个文档 3 克的数量不是太大,这就是我要做的:

    1. 按 3-gram 对两个文件进行键控以获得 (3-gram, doc#) 元组。
    2. cogroup 两个 RDD,为您提供 3gram 密钥和 2 个文档列表#
    3. 在单个 scala 函数中处理这些,输出一组(文档对)的所有唯一排列。
    4. 然后执行 coutByKeycountByKeyAprox 以计算每个文档对的不同 3-gram 的数量。

    注意:您可以跳过 .distinct() 与此通话。此外,您不应将每一行拆分两次。将line =&gt; (line.split("\t")(3),line.split("\t")(1))) 更改为line =&gt; { val s = line.split("\t"); (s(3),s(1)))

    编辑 2:

    你的记忆力似乎也调得很差。例如,使用.set("spark.shuffle.memoryFraction", "0.4").set("spark.storage.memoryFraction", "0.6") 基本上没有用于任务执行的内存(因为它们加起来为 1.0)。我应该早点看到这一点,但我专注于问题本身。

    查看调整指南 herehere

    另外,如果您在单台机器上运行它,您可能会尝试使用一个巨大的执行器(甚至完全放弃 Spark),因为您不需要分布式处理平台的开销(以及分布式硬件容错, ETC)。

    【讨论】:

    • 谢谢丹尼尔。这是我尝试加入的第一件事。我编辑了我的问题。但它又失败了。我是否错误地使用了连接功能?第一次加入时,我使用 3-gram 作为键,我有 40,000 个不同的 3-gram。现在我正在尝试 countbykey 定义。我不确定它是否会有所帮助:)
    • 嗯,让我看看我是否正确理解了这个问题。您可能会使用效率更高的aggregateByKey100772C111doc1是什么关系?此列在所有行上是否具有相同的值?如果不是,那么,您不是更愿意通过第二列值的每个排列进行汇总吗? (即100772C111-200772C222)?
    • 是的,很抱歉缺少信息。 100772C111 是文档的身份号码。 text1 中有 1.7 万个不同的 doc1,text2 中有 200 万个。并且在这两个文本中都没有共同的文档编号。对于每个文档编号,有 3 克属于它们。是的,我尝试汇总每个不同文档对的常见 3-gram 数量,一个来自 text1,另一个来自 text2。我编辑了文本示例。
    • 每个文件有 2 列(3gram 和 doc#)还是 4?我不知道你为什么有split("\t")(3),索引似乎太大了。你真的能把 File1 和 File2 的内容分开吗?
    • 文本文件中还有一些其他不相关的列。对于第一个,我希望出现在索引 0 和 5 上。对于其他 1 和 3。映射和拆分这些文件后,我得到 3-gram,doc#
    猜你喜欢
    • 2015-09-25
    • 1970-01-01
    • 1970-01-01
    • 2019-07-22
    • 2020-03-06
    • 1970-01-01
    • 2019-07-02
    • 1970-01-01
    • 2012-09-16
    相关资源
    最近更新 更多