【问题标题】:Is it possible to join two rdds' values to avoid expensive shuffling?是否可以加入两个 rdds 的值以避免昂贵的洗牌?
【发布时间】:2016-01-15 14:52:10
【问题描述】:

我有两个 RDD,它们都有两列作为 (K,V)。在这些 RDD 的源中,键出现在另一个之下,并且对于每一行,为键分配了不同且不同的值。创建 RDD 的文本文件在这篇文章的底部给出。

两个 RDD 中的键完全不同,我想根据它们的值加入两个 RDD,并尝试找出每对存在多少个共同值。例如我试图达到诸如 (1-5, 10) 之类的结果,这意味着来自 RDD1 的键值“1”和来自 RDD2 的键值“5”共享 10 个值。

我在一台具有 256 GB 内存和 72 个内核的机器上工作。一个文本文件为 500 MB,而另一个为 3 MB。

这是我的代码:

val conf = new SparkConf().setAppName("app").setMaster("local[*]").set("spark.shuffle.spill", "true")
.set("spark.shuffle.memoryFraction", "0.4")
.set("spark.executor.memory","128g")
.set("spark.driver.maxResultSize", "0")

val RDD1 = sc.textFile("\\t1.txt",1000).map{line => val s = line.split("\t"); (s(0),s(1))}

val RDD2 = sc.textFile("\\t2.txt",1000).map{line => val s = line.split("\t"); (s(1),s(0))}


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

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

    joined.foreach(println)

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

我尝试通过使用从脚本中看到的广播变量来管理这个问题。如果我加入 RDD2(有 250K 行),其自身对显示在相同的分区中,因此会发生更少的随机播放,因此需要 3 分钟才能获得结果。然而,当应用 RDD1 与 RDD2 时,配对分散在分区中,导致非常昂贵的洗牌过程,它总是最终给出

ERROR TaskSchedulerImpl: Lost executor driver on localhost: Executor heartbeat timed out after 168591 ms error.

根据我的结果:

  • 我是否应该尝试对文本文件进行分区以在较小的块中创建 RDD1 并用 RDD2 分别加入那些较小的块?

  • 还有其他方法可以根据它们的 Value 字段加入两个 RDD 吗?如果我将原始值描述为键并使用连接函数将它们连接起来,则值对再次分散在分区上,这再次导致非常昂贵的 reducebykey 操作。例如

    val RDD1 = sc.textFile("\\t1.txt",1000).map{line => val s = line.split("\t"); (s(1),s(0))}
    
        val RDD2 = sc.textFile("\\t2.txt",1000).map{line => val s = line.split("\t"); (s(1),s(0))}
    

    RDD1.join(RDD2).map(line => (line._2,1)).reduceByKey((a,b) => (a+b))

伪数据样本:

KEY VALUE
1   13894
1   17376
1   15688
1   22434
1   2282
1   14970
1   11549
1   26027
1   2895
1   15052
1   20815
2   9782
2   3393
2   11783
2   22737
2   12102
2   10947
2   24343
2   28620
2   2486
2   249
2   3271
2   30963
2   30532
2   2895
2   13894
2   874
2   2021
3   6720
3   3402
3   25894
3   1290
3   21395
3   21137
3   18739
...

一个小例子

RDD1

2   1
2   2
2   3
2   4
2   5
2   6
3   1
3   6
3   7
3   8
3   9
4   3
4   4
4   5
4   6

RDD2

21  1
21  2
21  5
21  11
21  12
21  10
22  7
22  8
22  13
22  9
22  11

基于此数据连接结果:

(3-22,1)
(2-21,1)
(3-22,1)
(2-21,1)
(3-22,1)
(4-21,1)
(2-21,1)
(3-21,1)
(3-22,1)
(3-22,1)
(2-21,1)
(3-22,1)
(2-21,1)
(4-21,1)
(2-21,1)
(3-21,1)

减少关键结果:

(4-21,1)
(3-21,1)
(2-21,3)
(3-22,3)

【问题讨论】:

  • 你应该展示一个使用 nano-data 的例子并告诉我 spark-default.conf 的内容是什么
  • 我刚刚编辑了问题。
  • 我在哪里工作 我们在一个集群中有 4 台 8GB 的​​计算机,我们能够毫无问题地读取 4GB 以上的文件,所以我敢打赌你的 $SPARK_HOME/conf/spark-defaults 有问题.conf,再加上你没有添加我要求的示例与 IN nano-data -> OUT nano-data(为了让我们更清楚)
  • 我添加了一个小例子。是的 Alberto,我认为我的来源必须足以处理这些数据。我正在通过spark.apache.org/docs/latest/tuning.html 阅读有关调整火花的信息。我在那里遇到数据局部性问题。所以我的文本文件存在于不同的驱动程序中。我的意思是数据在 S 驱动程序中,而我的火花在 D 驱动程序中。这有关系吗?谢谢。

标签: scala join apache-spark


【解决方案1】:

您是否考虑过使用笛卡尔连接?您也许可以尝试以下方法:

val rdd1 = sc.parallelize(for { x <- 1 to 3; y <- 1 to 5 } yield (x, y)) // sample RDD
val rdd2 = sc.parallelize(for { x <- 1 to 3; y <- 3 to 7 } yield (x, y)) // sample RDD with slightly displaced values from the first

val g1 = rdd1.groupByKey()
val g2 = rdd2.groupByKey()

val cart = g1.cartesian(g2).map { case ((key1, values1), (key2, values2)) => 
             ((key1, key2), (values1.toSet & values2.toSet).size) 
           }

当我尝试在集群中运行上述示例时,我看到以下内容:

scala> rdd1.take(5).foreach(println)
...
(1,1)
(1,2)
(1,3)
(1,4)
(1,5)
scala> rdd2.take(5).foreach(println)
...
(1,3)
(1,4)
(1,5)
(1,6)
(1,7)
scala> cart.take(5).foreach(println)
...
((1,1),3)
((1,2),3)
((1,3),3)
((2,1),3)
((2,2),3)

结果表明对于(key1,key2),集合之间有3个匹配元素。请注意,这里的结果始终为 3,因为初始化的输入元组的范围与 3 个元素重叠。

笛卡尔变换也不会导致洗牌,因为它只是迭代每个 RDD 的元素并产生笛卡尔积。您可以通过在示例中调用 toDebugString() 函数来看到这一点。

scala> val carts = rdd1.cartesian(rdd2)
carts: org.apache.spark.rdd.RDD[((Int, Int), (Int, Int))] = CartesianRDD[9] at cartesian at <console>:25

scala> carts.toDebugString
res11: String =
(64) CartesianRDD[9] at cartesian at <console>:25 []
 |   ParallelCollectionRDD[1] at parallelize at <console>:21 []
 |   ParallelCollectionRDD[2] at parallelize at <console>:21 []

【讨论】:

  • 抱歉,这只会让情况变得更糟。
  • @zero323 哎呀,是的,我没有考虑内存后果。如果较小的数据集有 250K 条记录并且数据集大致相似,那么较大的数据集有大约 41M 条记录。完整的笛卡尔达到大约 10TB 的大小,这是非常不可行的。
  • 不仅如此。即使您只考虑网络流量,cartesian 也会更糟。加入一对一的映射,但笛卡尔积映射一对一。这意味着每个分区都必须移动多次。
猜你喜欢
  • 2019-04-28
  • 1970-01-01
  • 1970-01-01
  • 2021-02-01
  • 2012-12-30
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
  • 1970-01-01
相关资源
最近更新 更多